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..0651288 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,14 +175,21 @@ ;; 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]))) +;; Leaving also drops the socket. It isn't just tidiness any more: an open +;; socket keeps us in the project's roster, so peers would go on seeing us +;; viewing something we walked away from — and its deltas would land in the +;; next project's db while that one loads. (rf/reg-event-fx ::nav-list (fn [{:keys [db]} _] {:db (-> (reset-project db) (assoc :page :list)) + :connect-scene nil :fetch-projects true})) (rf/reg-event-db ::set-projects (fn [db [_ ps]] (assoc db :projects ps :projects-error nil))) -(rf/reg-event-db ::nav-create (fn [db _] (-> (reset-project db) - (assoc :page :create :create-error nil)))) +(rf/reg-event-fx ::nav-create (fn [{:keys [db]} _] + {:db (-> (reset-project db) + (assoc :page :create :create-error nil)) + :connect-scene nil})) ;; loading a project is a three-hop chain: detail → OTIO file → scene, each ;; feeding the next, and any hop's failure lands the editor in :error. @@ -148,6 +199,7 @@ (assoc :page :editor) (assoc-in [:load :status] :loading) (assoc-in [:view :route-state] view-state)) + :connect-scene nil ; the previous project's, until ::project-ready :http-xhrio (api/GET (str "/api/projects/" id "/") {:on-success [::project-detail-loaded] :on-failure [::project-load-error]})})) @@ -190,16 +242,230 @@ ;; 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)))))) + +;; Changing party drops anything parked: it's a position from the party we just +;; left, and replaying it later (on the way out of a form, say) would drag us +;; back to people we're no longer with. +(defn- set-party [db pid] + (-> db (assoc-in [:peers :roster (my-cid db) :party] pid) + (assoc-in [:peers :parked] nil))) + +(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 (set-party db (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 (my-cid db) + (let [db (set-party db 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?)}))) + +;; Nav is the only message that steers another browser, and it arrives relayed +;; from whatever a peer chose to send. Nothing downstream should have to cope +;; with a shape it can't use — a bad frame reaching scene/assert-frame would +;; throw inside the event loop and take the editor down with it. +(defn- nav-stack + "The peer's stack as context keywords, or nil if it isn't one we could stand + in: it has to start at the root and name timelines all the way down." + [stack] + (let [ks (when (sequential? stack) (mapv #(when (string? %) (keyword %)) stack))] + (when (and (seq ks) (= :root (first ks)) (every? some? ks)) ks))) + +(defn- nav-frame [n] + (when (and (number? n) (not (neg? n)) (== n (js/Math.floor n))) n)) + +(defn- standable? [db gid] + (contains? #{:timeline :annotation} (get-in db [:scene :groups gid :type]))) + +;; 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 playing] :as msg}]] + (let [stack (nav-stack (:stack msg)) + playhead (nav-frame (:playhead msg))] + (cond + (not (party/mirrors? (roster db) (my-cid db) cid)) + ;; not (or no longer) one of ours — and if we were holding their move + ;; from when they were, it's stale now + {:db (cond-> db (= cid (get-in db [:peers :parked :cid])) + (assoc-in [:peers :parked] nil))} + + (not (and stack playhead)) {:db db} ; nothing we can act on + + ;; 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 (every? #(standable? db %) 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] 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 +877,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