1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673 | ;; ported-from: src/lib/media.ts — and, since Phase 6b, no longer a
;; transliteration of it.
;;
;; Media-only WebRTC for room calls (PROTOCOL.md v2 §7c). libp2p's WebRTC
;; transport carries data only, so calls run their own per-pair
;; RTCPeerConnections carrying nothing but audio/video (DTLS-SRTP, own
;; ICE/TURN). SDP/ICE ride sealed signaling envelopes (kSig — the same keys and
;; AADs as v1 signaling) over the /sueta/2/call-sig/1.0 libp2p stream, so an
;; established call keeps flowing even if the libp2p path dies — signaling is
;; only needed for (re)negotiation.
;;
;; The negotiation machinery — perfect negotiation with politeness, transceiver
;; reuse, replaceTrack mutes, restartIce + deterministic rebuild — is lifted from
;; the v1 mesh minus everything datachannel.
;;
;; ── THE TWO MUTEXES, AND WHY THEY ARE NOT `p/let` ────────────────────────────
;;
;; Every peer record carries two promise chains used as LOCKS:
;;
;; `queue` serialises SIGNAL handling. Two offers arriving back to back
;; must be applied one after the other; a second
;; setRemoteDescription entered while the first is still
;; in-flight is a WebRTC state-machine violation.
;; `mediaQueue` serialises MEDIA changes, so an attach and a detach racing on
;; the same transceiver cannot both decide what it holds.
;;
;; `fx/locked!` is that pattern with a name. It matters that it has one:
;; `p/let` preserves ordering WITHIN one handler and destroys it BETWEEN two,
;; so replacing either chain with a `p/let` reads exactly like a simplification
;; and silently interleaves the handlers. test/vectors/media-effects.json's
;; `signalMutex` scenario exists for that substitution and nothing else — two
;; frames delivered without awaiting either, and A's whole sequence must appear
;; before B's first step.
;;
;; ── WHAT ELSE IS PINNED, AND HOW EXACTLY ─────────────────────────────────────
;;
;; The five scenarios in media-effects.json are microtask-precise transcripts:
;; they record the interleaving of the two `apply-media!` runs a first attach
;; makes (see `attach-media`'s note) — serialised on `mediaQueue` since the
;; duplicate-addTrack fix, but still two — so the promise composition of
;; `apply-media!` and of `setup-pc!`'s tail is pinned to the hop, not just to
;; the order. Both are therefore written as explicit `.then` chains rather than
;; `p/let`: what looks here like a leftover from the transliteration is a
;; deliberate refusal to move a fixture. Everything strictly sequential —
;; `handle-signal` and below — IS `p/let`, and the transcripts prove it.
;;
;; A SignalPayload is {kind:"offer"|"answer"|"ice", sdp?|candidate?, name} —
;; small, but JSON-serialised into a kSig seal, so it is built with `j/ordered`
;; and NOT through sueta.wire/encode: `encode` omits a nil field, and the
;; end-of-candidates signal is exactly `{kind:"ice", candidate:null, name}`,
;; whose null is on the wire that deployed peers already speak.
(ns sueta.lib.media
(:require [promesa.core :as p]
[sueta.fx :as fx]
[ardegazu.rooms.js :as j]
[ardegazu.rooms.lib.crypto :as rc]
[shadow.cljs.modern :refer (defclass)]))
(def CALL-SIG-PROTOCOL "/sueta/2/call-sig/1.0")
(def ^:private FAILED-REBUILD-MS 10000)
(def ^:private KINDS ["audio" "video"])
(def ^:private te (js/TextEncoder.))
(def ^:private td (js/TextDecoder.))
;; ---- the peer record -------------------------------------------------------
;;
;; A plain mutable JS object, deliberately: `polite`, `ignoreOffer`,
;; `settingRemoteAnswer` and `pendingCandidates` are read straight off it by
;; test/helpers/media-effects.mjs through the `_peer` seam, so these field names
;; are part of a test surface. Every read of one goes through the accessors
;; below, which is what keeps the negotiation logic itself free of string keys.
(defn- mk-media-peer
"`polite` is deterministic and asymmetric per pair — the smaller base58 PeerId
is polite — which is what makes a glare collision resolvable without a round
trip. CLJS's two-argument `<` inlines to the JS operator, so this is a string
comparison, not a numeric one."
[id my-id]
(j/ordered "id" id
"polite" (< my-id id)
"pc" nil
"makingOffer" false
"ignoreOffer" false
"settingRemoteAnswer" false
"pendingCandidates" (array)
"wantedStream" nil
"maxVideoKbps" js/undefined
"mediaQueue" (js/Promise.resolve nil)
"queue" (js/Promise.resolve nil)
"rebuildTimer" nil))
(defn- peer-id [peer] (unchecked-get peer "id"))
(defn- ^js peer-pc [peer] (unchecked-get peer "pc"))
(defn- wanted-stream [peer] (unchecked-get peer "wantedStream"))
(defn- max-kbps [peer] (unchecked-get peer "maxVideoKbps"))
(defn- pending-candidates [peer] (unchecked-get peer "pendingCandidates"))
(defn- ^boolean polite? [peer] (j/truthy? (unchecked-get peer "polite")))
(defn- ^boolean making-offer? [peer] (j/truthy? (unchecked-get peer "makingOffer")))
(defn- ^boolean ignoring-offer? [peer] (j/truthy? (unchecked-get peer "ignoreOffer")))
(defn- ^boolean answering? [peer] (j/truthy? (unchecked-get peer "settingRemoteAnswer")))
;; a timer id of 0 is a real id, and 0 is truthy in Clojure — hence j/truthy?
(defn- ^boolean rebuild-armed? [peer] (j/truthy? (unchecked-get peer "rebuildTimer")))
(defn- clear-rebuild! [peer]
(when (rebuild-armed? peer) (js/clearTimeout (unchecked-get peer "rebuildTimer")))
(unchecked-set peer "rebuildTimer" nil)
js/undefined)
(defn- destroy-peer! [peer]
(when (rebuild-armed? peer)
(js/clearTimeout (unchecked-get peer "rebuildTimer")))
;; closing a pc that is already closed throws in some hosts
(try
(when-some [pc (peer-pc peer)] (.close pc))
(catch :default _ nil))
js/undefined)
(defn- ^boolean dead?
"A pc that will never carry media again — either gone or terminally failed."
[pc]
(or (not (j/truthy? pc))
(identical? "failed" (.-connectionState ^js pc))
(identical? "closed" (.-connectionState ^js pc))))
(defn- ^boolean track-of-kind?
"TS `tr.receiver.track?.kind === kind`. A transceiver has a receiver track
from the moment it exists, muted or not, which is how an existing one is
found for reuse instead of adding a second."
[tx kind]
(let [track (.-track (.-receiver ^js tx))]
(identical? kind (when (j/truthy? track) (.-kind ^js track)))))
(declare setup-pc! apply-media! send-signal! on-sig-frame handle-signal
add-candidate! classify-transport! rebuild! attach-media)
;; ---- the manager -----------------------------------------------------------
(defclass Media
(constructor [this opts]
(unchecked-set this "_crypto" (unchecked-get opts "crypto"))
(unchecked-set this "_ice" (unchecked-get opts "ice"))
(unchecked-set this "_net" (unchecked-get opts "net"))
(unchecked-set this "_myName" (unchecked-get opts "myName"))
(unchecked-set this "_ev" (unchecked-get opts "events"))
(unchecked-set this "_peers" (js/Map.))
(unchecked-set this "_closed" false)
(.handleFrames ^js (unchecked-get opts "net") CALL-SIG-PROTOCOL
(fn [from frame] (on-sig-frame this from frame) js/undefined))))
;; The JS boundary for the manager: seven fields, named once each.
(defn- ^js peers-of [self] (unchecked-get self "_peers"))
(defn- ^js net-of [self] (unchecked-get self "_net"))
(defn- my-id [self] (unchecked-get (net-of self) "myId"))
(defn- my-name [self] ((unchecked-get self "_myName")))
(defn- ^boolean closed? [self] (j/truthy? (unchecked-get self "_closed")))
(defn- emit! [self ev & args]
(apply (unchecked-get (unchecked-get self "_ev") ev) args)
js/undefined)
(defn attach-media
"Declaratively set the media sent to one peer; creates the media pc on first
use. Diffs against current senders; a stream missing a kind soft-mutes it.
Survives pc rebuilds.
TWO THINGS HERE ARE LOAD-BEARING AND LOOK LIKE ACCIDENTS.
`setup-pc!` is NOT awaited, and `wantedStream` is written immediately after
it — that ordering is the whole rebuild story: by the time the ICE config
resolves and the pc exists, the field the pc's tail reads is already set.
On a FIRST attach the media is still applied twice — once by `setup-pc!`'s
tail, once by the `mediaQueue` chain below — but the two are SERIALISED, because
that tail chains onto `mediaQueue` rather than calling `apply-media!` beside
it. The later run therefore sees the transceiver the earlier one added and
takes the reuse path.
It did not always. Until that tail was routed through the mutex the pair ran
concurrently, both read `getTransceivers()` before either `addTrack` landed,
both added the same track, and the browser rejected the second with
InvalidAccessError — swallowed by `apply-media!`'s own catch into a
console.warn. Harmless in effect (one transceiver, one sender, the right
track) but real, and it was pinned by media-effects.json and by a test that
named it out loud, precisely so that fixing it could not be silent. Fixing it
removed four rows from each of the five transcripts and nothing else: every
glare, politeness, ICE-buffering and mutex property is byte-identical across
the change."
[self id stream opts]
(when-not (closed? self)
(let [peers (peers-of self)
existing (.get peers id)
peer (if (identical? existing js/undefined)
(let [fresh (mk-media-peer id (my-id self))]
(.set peers id fresh)
(setup-pc! self fresh)
fresh)
existing)]
(unchecked-set peer "wantedStream" stream)
;; TS `if (opts?.maxVideoKbps !== undefined)` — an absent option must not
;; erase a cap a previous attach set
(let [kbps (when (j/truthy? opts) (unchecked-get opts "maxVideoKbps"))]
(when-not (identical? kbps js/undefined) (unchecked-set peer "maxVideoKbps" kbps)))
(fx/locked! peer "mediaQueue" #(apply-media! self peer))))
js/undefined)
(defn- deactivate!
"One transceiver, sender-less and inactive, in one tick. A pc closed
underneath us rejects the direction write."
[tx]
(try (set! (.-direction ^js tx) "inactive") (catch :default _ nil))
js/undefined)
(defn- stop-receiving!
"Every transceiver goes sender-less and inactive. Sequential: two
`replaceTrack` calls on one pc may not overlap.
NO GUARD HERE, AND THAT IS CORRECT — the story, because it reads like an
omission. The TypeScript opened its loop with
`if (t.receiver.track?.kind === undefined) continue;`, and the port rendered
that as `(identical? js/undefined (when (truthy? track) (.-kind track)))`,
which is `null === undefined` and so always false: the skip never fired in
any shipped CLJS build.
That was written up as a latent behaviour difference until it was actually
measured, and the measurement says otherwise: `receiver.track` is never
null. Checked in Chromium across 18 transceivers in every shape these pcs
hold — created by `addTrack` before negotiation, explicit `sendonly`,
explicit `recvonly`, after a full negotiation with a peer that sends nothing
back, on the answering side, and after `replaceTrack(null)` +
`direction=\"inactive\"` — the expression was false every time. Per spec an
`RTCRtpReceiver` always carries a MediaStreamTrack, so the `?.` guarded a
state WebRTC does not produce.
The guard was therefore dead in the ORIGINAL too, this loop matches the
TypeScript's real behaviour, and restoring the skip would add a branch no
conforming browser can reach. Do not \"fix\" it — and in particular not to
`some?`, which turns an unreachable branch into a throw the moment a
receiver is missing."
[pc]
(j/each-in-order!
(array-seq (js/Array.from (.getTransceivers ^js pc)))
(fn [tx]
(p/let [_ (fx/drop-silently (.replaceTrack (.-sender ^js tx) nil))]
(deactivate! tx)))))
(defn detach-media
"Stop sending AND receiving media with this peer, without tearing the pc down
— the pair may still be in the room."
[self id]
(let [peer (.get ^js (peers-of self) id)]
(when (and (j/truthy? peer) (j/truthy? (peer-pc peer)))
(unchecked-set peer "wantedStream" nil)
(fx/locked!
peer "mediaQueue"
(fn []
;; re-attached while we waited for the lock — skip the stale detach
(if (j/truthy? (wanted-stream peer))
(p/resolved nil)
(stop-receiving! (peer-pc peer)))))))
js/undefined)
(defn detach-all-media [self]
(doseq [id (array-seq (js/Array.from (.keys ^js (peers-of self))))]
(detach-media self id))
js/undefined)
(defn drop-peer
"Tear the pair's media pc down entirely (peer left the call/room)."
[self id]
(let [peers (peers-of self)
peer (.get peers id)]
(when-not (identical? peer js/undefined)
(destroy-peer! peer)
(.delete peers id)))
js/undefined)
(defn close [self]
(unchecked-set self "_closed" true)
(.forEach ^js (peers-of self) (fn [p _k] (destroy-peer! p)))
(.clear ^js (peers-of self))
js/undefined)
(defn debug-state
"One row per peer for the e2e harness: who, what state, which kinds we send."
[self]
(let [out (array)]
(.forEach ^js (peers-of self)
(fn [p _k]
(let [pc (peer-pc p)]
(.push out (j/ordered "id" (peer-id p)
"conn" (when (j/truthy? pc) (.-connectionState pc))
"sending" (when (j/truthy? pc)
(.map (.filter (js/Array.from (.getSenders pc))
(fn [s] (j/truthy? (.-track ^js s))))
(fn [s] (.-kind (.-track ^js s))))))))))
out))
;; ---- pc lifecycle (lifted from the v1 mesh) --------------------------------
(defn- on-negotiation-needed!
"Our own offer. `makingOffer` brackets it, and it is what the far side's
glare test reads about us."
[self peer pc]
(-> (j/later
(fn []
;; the flag goes up in the SAME tick the offer starts, which is what
;; the far side's glare test is reading about us
(unchecked-set peer "makingOffer" true)
(p/let [_ (.setLocalDescription ^js pc)]
(send-signal! self (peer-id peer)
(j/ordered "kind" "offer"
"sdp" (.-sdp (.-localDescription ^js pc))
"name" (my-name self))))))
(.catch (fn [err] (js/console.warn "call negotiation failed" err) nil))
(.then (fn [_] (unchecked-set peer "makingOffer" false) js/undefined)))
js/undefined)
(defn- on-connection-failed!
"ICE gave up. Restart it, and arm ONE deterministic rebuild: the pair's
smaller id is the single re-initiator, so a failure never produces two
simultaneous rebuilds. TS wrote `peer.rebuildTimer ??= setTimeout(…)`, which
is NULLISH — an already-armed timer is not re-armed."
[self peer pc]
(.restartIce ^js pc)
(when-not (rebuild-armed? peer)
(unchecked-set
peer "rebuildTimer"
(js/setTimeout (fn []
(unchecked-set peer "rebuildTimer" nil)
(when (and (dead? pc) (< (my-id self) (peer-id peer)))
(rebuild! self peer))
js/undefined)
FAILED-REBUILD-MS)))
js/undefined)
(defn- install-handlers!
"The four events a media-only pc raises. Assigned rather than added, exactly
as the v1 mesh did: one handler each, and a rebuilt pc starts clean."
[self peer pc]
(set! (.-ontrack pc)
(fn [^js e]
;; a track with no stream of its own still has to be shown
(let [s (aget (.-streams e) 0)]
(emit! self "track" (peer-id peer)
(if (identical? s js/undefined) (js/MediaStream. #js [(.-track e)]) s)))
js/undefined))
(set! (.-onnegotiationneeded pc) (fn [] (on-negotiation-needed! self peer pc)))
(set! (.-onicecandidate pc)
(fn [^js e]
(send-signal! self (peer-id peer)
;; a null candidate is END-OF-CANDIDATES and is sent as
;; an explicit null — deployed peers read the key
(j/ordered "kind" "ice"
"candidate" (if (j/truthy? (.-candidate e))
(.toJSON (.-candidate e))
nil)
"name" (my-name self)))
js/undefined))
(set! (.-onconnectionstatechange pc)
(fn []
(case (.-connectionState pc)
"connected" (do (clear-rebuild! peer) (classify-transport! self peer))
"failed" (on-connection-failed! self peer pc)
nil)
js/undefined))
js/undefined)
(defn- setup-pc!
"Build the pair's pc once the ICE config is known, then give it whatever media
the app already wants — that last line is the rebuild path.
An explicit `.then`, not `p/let`: this is one await, and its microtask
position relative to `attach-media`'s `mediaQueue` chain is recorded in every
media-effects transcript."
[self peer]
(.then
^js (.get ^js (unchecked-get self "_ice"))
(fn [config]
;; closed, or a concurrent caller won the race to build it
(when-not (or (closed? self) (j/truthy? (peer-pc peer)))
(let [pc (js/RTCPeerConnection. config)]
(unchecked-set peer "pc" pc)
(install-handlers! self peer pc)
;; THROUGH the mutex, not beside it. `attach-media` chains its own
;; `apply-media!` onto `mediaQueue`; this tail used to call it
;; directly, so on a first attach the two ran concurrently, both read
;; `getTransceivers()` before either `addTrack` landed, and the
;; browser rejected the second with InvalidAccessError. Chaining it
;; here serialises the pair: the later run sees the transceiver the
;; earlier one added and takes the reuse path.
(when (j/truthy? (wanted-stream peer))
(fx/locked! peer "mediaQueue" #(apply-media! self peer)))
js/undefined)))))
(defn- wanted-track
"The track of `kind` the app wants sent, or nil. TS `?? null`."
[stream kind]
(when (j/truthy? stream)
(let [t (.find (js/Array.from (.getTracks ^js stream))
(fn [tr] (identical? (.-kind ^js tr) kind)))]
(when-not (identical? t js/undefined) t))))
(defn- apply-kind!
"Reconcile ONE kind against the pc: reuse the transceiver if there is one,
add it if there is not, and soft-mute (replaceTrack nil, no renegotiation) if
the app no longer wants to send this kind."
[pc want tx stream]
(cond
(j/truthy? want)
(if-not (identical? tx js/undefined)
(p/let [_ (if (identical? (.-track (.-sender ^js tx)) want)
nil
(.replaceTrack (.-sender ^js tx) want))]
(when-not (identical? "sendrecv" (.-direction ^js tx))
(set! (.-direction ^js tx) "sendrecv"))
js/undefined)
;; no transceiver for this kind yet — addTrack makes one
(do (.addTrack ^js pc want stream) js/undefined))
(and (not (identical? tx js/undefined)) (j/truthy? (.-track (.-sender ^js tx))))
(.replaceTrack (.-sender ^js tx) nil)
:else js/undefined))
(defn- apply-bitrate-cap!
"A relayed pair can be given a lower ceiling than a direct one. Best effort:
`setParameters` rejects on a pc that moved underneath us."
[pc kbps]
(let [sender (.find (js/Array.from (.getSenders ^js pc))
(fn [s] (let [track (.-track ^js s)]
(identical? "video" (when (j/truthy? track) (.-kind ^js track))))))]
(if (or (identical? sender js/undefined) (not (j/truthy? kbps)))
(js/Promise.resolve nil)
(fx/drop-silently
(j/later
(fn []
(let [params (.getParameters ^js sender)
enc (.-encodings params)]
(set! (.-encodings params)
(if (and (j/truthy? enc) (j/truthy? (.-length ^js enc))) enc #js [(js-obj)]))
(unchecked-set (aget (.-encodings params) 0) "maxBitrate" (* kbps 1000))
(.setParameters ^js sender params))))))))
;; `_self` — every helper in this family takes the Media receiver first and the
;; call sites pass it; this one happens to read nothing off it. Kept in place so
;; the family keeps one shape.
(defn- apply-media!
"Make the pc send exactly what `wantedStream` holds, audio then video, then
apply the bitrate cap.
Explicit `.then` links rather than `p/let`, for the reason the namespace
header gives: two of these still run on a first attach — serialised on
`mediaQueue`, not concurrent — and every media-effects transcript records the
resulting interleaving to the microtask, so a `p/let` here moves a fixture.
The per-kind `.catch` no longer has a duplicate to swallow, but it stays: a
pc that moves underneath a kind still rejects, and one kind failing must not
abandon the other."
[_self peer]
(let [pc (peer-pc peer)]
(if (or (not (j/truthy? pc)) (identical? "closed" (.-connectionState pc)))
(js/Promise.resolve nil)
(-> ((fn kind-step [i]
(if (>= i (count KINDS))
(js/Promise.resolve nil)
(let [kind (nth KINDS i)
stream (wanted-stream peer)
want (wanted-track stream kind)
tx (.find (js/Array.from (.getTransceivers pc))
(fn [tr] (track-of-kind? tr kind)))]
(-> (js/Promise.resolve nil)
(.then (fn [_] (apply-kind! pc want tx stream)))
(.catch (fn [err] (js/console.warn "applyMedia" kind err) nil))
(.then (fn [_] (kind-step (inc i))))))))
0)
(.then (fn [_] (apply-bitrate-cap! pc (max-kbps peer))))))))
;; ---- sealed signaling over the libp2p stream --------------------------------
(defn- send-signal!
"Seal one SignalPayload to `to` and put it on the call-sig stream. Fire and
forget: a peer that left mid-send is expected, not news."
[self to payload]
(fx/drop-silently
(p/let [env (rc/seal-signal (unchecked-get self "_crypto") to payload)]
(.sendFrame ^js (net-of self) CALL-SIG-PROTOCOL to
(.encode te (js/JSON.stringify env)))))
js/undefined)
(defn- peer-for
"The peer record for `id`, created on first sight — on the callee side the
offerer's first signal is what brings the pair into existence, and its pc
must exist before the signal is applied, so this one DOES await `setup-pc!`."
[self id]
(let [peers (peers-of self)
existing (.get peers id)]
(if (identical? existing js/undefined)
(let [fresh (mk-media-peer id (my-id self))]
(.set peers id fresh)
(p/let [_ (setup-pc! self fresh)] fresh))
existing)))
(defn- on-sig-frame
"One sealed signaling frame. Anything that will not parse or will not open is
dropped in silence (PROTOCOL.md §3): a frame from a non-member is
indistinguishable from noise, and saying so would be a side channel."
[self from frame]
(j/later
(fn []
(let [env (j/parse-json (.decode td frame))]
(when (some? env)
(p/let [payload (rc/open-signal (unchecked-get self "_crypto") from env)]
(when (some? payload)
(p/let [peer (peer-for self from)]
;; THE MUTEX — see the namespace header
(do (fx/locked! peer "queue" #(handle-signal self peer payload))
js/undefined))))))))
js/undefined)
(defn- flush-candidates!
"The remote description has landed. Close the answer window, then drain the
buffer in ARRIVAL order — ICE is a priority-ordered list. Spliced, not
copied, so a second flush cannot replay it."
[peer]
(unchecked-set peer "settingRemoteAnswer" false)
(let [pending (.splice (pending-candidates peer) 0)]
(j/each-in-order! (array-seq pending) #(add-candidate! peer %))))
(defn- ^boolean ready-for-offer?
"Perfect negotiation's readiness test. The second clause is the whole reason
an impolite peer can still accept an offer that arrives in
\"have-local-offer\": the ANSWER it is in the middle of applying is about to
make the state stable anyway. Close that window and the identical offer is
dropped instead."
[peer pc]
(and (not (making-offer? peer))
(or (identical? "stable" (.-signalingState ^js pc))
(answering? peer))))
(defn- apply-description!
"An offer or an answer. A remote OFFER arriving when we are not ready is a
collision, and politeness decides it: the polite side rolls back and accepts,
the impolite side sets `ignoreOffer` and never tells the pc about it at all."
[self peer payload kind]
(let [pc (peer-pc peer)
description (j/ordered "type" kind "sdp" (unchecked-get payload "sdp"))
collision? (and (identical? "offer" kind) (not (ready-for-offer? peer pc)))]
(unchecked-set peer "ignoreOffer" (and (not (polite? peer)) collision?))
(when-not (ignoring-offer? peer)
(unchecked-set peer "settingRemoteAnswer" (identical? "answer" kind))
(p/let [_ (.setRemoteDescription pc description)
_ (flush-candidates! peer)]
(when (identical? "offer" kind)
(p/let [_ (.setLocalDescription pc)]
(send-signal! self (peer-id peer)
(j/ordered "kind" "answer"
"sdp" (.-sdp (.-localDescription pc))
"name" (my-name self)))))))))
(defn- apply-ice!
"A candidate arriving before any remote description would throw in the
browser, so it waits on `pendingCandidates` until one lands."
[peer candidate]
(if-not (j/truthy? (.-remoteDescription (peer-pc peer)))
(do (.push (pending-candidates peer) candidate) js/undefined)
(add-candidate! peer candidate)))
(defn- handle-signal
"One signal, applied under the peer's `queue` lock. Never rejects: a
negotiation fault must not poison the lock for every later signal."
[self peer payload]
(-> (j/later
(fn []
(let [kind (unchecked-get payload "kind")]
(case kind
("offer" "answer") (apply-description! self peer payload kind)
"ice" (apply-ice! peer (unchecked-get payload "candidate"))
js/undefined))))
(.catch (fn [err] (js/console.warn "call signal handling failed" err) nil))))
(defn- add-candidate!
"TS `candidate ?? undefined` — NULLISH, and `undefined` is how WebRTC is told
'end of candidates'. A rejection here is only worth a word when we are not
already deliberately ignoring this peer's offer."
[peer candidate]
(-> (j/later #(.addIceCandidate (peer-pc peer) (j/nn candidate js/undefined)))
(.catch (fn [err]
(when-not (ignoring-offer? peer)
(js/console.warn "addIceCandidate failed" err))
nil))))
;; ---- transport classification ----------------------------------------------
(defn- selected-local-candidate
"The local candidate of the nominated pair, folded out of one RTCStatsReport.
The transport report's own `selectedCandidatePairId` wins where the browser
fills it in; otherwise the nominated succeeded pair is taken. Pure over
`stats`."
[stats]
(let [locals (js/Map.)
selected (volatile! nil)]
(.forEach ^js stats
(fn [s]
(case (unchecked-get s "type")
"local-candidate" (.set locals (unchecked-get s "id") s)
"candidate-pair"
(when (and (j/truthy? (unchecked-get s "nominated"))
(identical? "succeeded" (unchecked-get s "state")))
(vreset! selected s))
"transport"
(when (j/truthy? (unchecked-get s "selectedCandidatePairId"))
(when-some [pair (.get ^js stats (unchecked-get s "selectedCandidatePairId"))]
(vreset! selected pair)))
nil)
js/undefined))
(when-some [sel @selected]
(.get locals (unchecked-get sel "localCandidateId")))))
(defn- classify-transport!
"Tell the UI whether this pair went direct or through a TURN relay. Any
failure reports \"direct\" — the honest default, since a relayed pair is the
one worth flagging."
[self peer]
(-> (p/let [stats (.getStats (peer-pc peer))]
(let [local (selected-local-candidate stats)]
(emit! self "mediaState" (peer-id peer)
(if (identical? "relay" (when (some? local) (unchecked-get local "candidateType")))
"relayed"
"direct"))))
(.catch (fn [_] (emit! self "mediaState" (peer-id peer) "direct") nil))))
(defn- rebuild!
"Throw the pc away and attach again — the deterministic recovery from a failed
pair, run by one side only (see `on-connection-failed!`)."
[self peer]
(let [stream (wanted-stream peer)
kbps (max-kbps peer)
id (peer-id peer)]
(drop-peer self id)
(when (j/truthy? stream)
(attach-media self id stream
(if (identical? kbps js/undefined) js/undefined (j/ordered "maxVideoKbps" kbps)))))
js/undefined)
;; ---- class surface ---------------------------------------------------------
(let [proto (.-prototype Media)]
(unchecked-set proto "attachMedia" (fn [id stream opts] (this-as self (attach-media self id stream opts))))
(unchecked-set proto "detachMedia" (fn [id] (this-as self (detach-media self id))))
(unchecked-set proto "detachAllMedia" (fn [] (this-as self (detach-all-media self))))
(unchecked-set proto "dropPeer" (fn [id] (this-as self (drop-peer self id))))
(unchecked-set proto "close" (fn [] (this-as self (close self))))
(unchecked-set proto "debugState" (fn [] (this-as self (debug-state self))))
;; the "private" surface the tests drive (peer-kit's `_name` idiom, the same
;; shape lib/log and lib/net expose). 445 lines of perfect-negotiation race
;; logic had no export, no seam and no test; the negotiation is where the
;; races are, and `_onSigFrame` is the only way in that also exercises the
;; per-peer promise-chain MUTEX, while `_handleSignal` is the only way to
;; reach the `settingRemoteAnswer` window (the queue closes it).
(unchecked-set proto "_peer" (fn [id] (this-as self (.get ^js (peers-of self) id))))
(unchecked-set proto "_setupPc" (fn [peer] (this-as self (setup-pc! self peer))))
(unchecked-set proto "_applyMedia" (fn [peer] (this-as self (apply-media! self peer))))
(unchecked-set proto "_sendSignal" (fn [to payload] (this-as self (send-signal! self to payload))))
(unchecked-set proto "_onSigFrame" (fn [from frame] (this-as self (on-sig-frame self from frame))))
(unchecked-set proto "_handleSignal" (fn [peer payload] (this-as self (handle-signal self peer payload))))
(unchecked-set proto "_addCandidate" (fn [peer candidate] (add-candidate! peer candidate)))
(unchecked-set proto "_classifyTransport" (fn [peer] (this-as self (classify-transport! self peer))))
(unchecked-set proto "_rebuild" (fn [peer] (this-as self (rebuild! self peer)))))
|