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
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945 | ;; ported-from: src/lib/net.ts @ v1.0.0
;;
;; libp2p transport for the sueta p2p stack (v2): one libp2p node per session
;; with an ephemeral Ed25519 identity, a websocket connection to a
;; circuit-relay-v2 node, peer discovery over the relay's shared pubsub topic,
;; and direct browser-to-browser WebRTC upgrades. See chat/docs/PROTOCOL.md v2.
;;
;; Membership model: a discovered peer is only surfaced to the app after it
;; proves room membership with a presence beacon that decrypts under the room's
;; keys (the v1 "undecryptable peers get dropped" rule, §3). The relay and
;; unrelated peers on it see only an opaque topic string and sealed envelopes;
;; room traffic itself never touches the relay — gossipsub rides the direct
;; member-to-member connections.
;;
;; SHAPE. Four bands, and the file reads top-down through them:
;; 1. the objects hanging off `self` — every `unchecked-get` chain, named once
;; 2. DECISIONS as plain values — redial-delay, conn-state, dial-target,
;; sweep-verdict take values and return values, so what the membership
;; machine decides is separable from what it then does about it
;; 3. the effectful shell, in dependency order
;; 4. `start`, split into one installer per subsystem, and the class surface
;; The file used to open with a 21-name `(declare)`, which reads like a fully
;; cyclic call graph. There is exactly ONE real cycle here — ensure-relay and
;; schedule-redial call each other, which is the retry loop — and it is the only
;; forward declaration left.
;;
;; TWO-STACK LAW (house rule 11): this is the libp2p **v2** side. @libp2p/webrtc
;; v5 here loads @ipshipyard/node-datachannel 0.26; peer-kit's v3 loads
;; node-datachannel 0.33, and two builds of libdatachannel in one process abort
;; it. Every import below stays an EXTERNAL ESM import (shadow-cljs.edn's
;; :js-provider :import) and deploy/check-dist.sh greps the emitted bundles.
;;
;; PORTING NOTE — JS truthiness is load-bearing throughout this file:
;; rec.lastBeacon 0 = never proved membership (a real timestamp otherwise)
;; #redialTimer 0/null = no redial scheduled
;; name undefined = unnamed
;; A bare `(when x …)` on any of those is always true in Clojure. Every such
;; test goes through ardegazu.rooms.js/truthy?.
(ns ardegazu.rooms.lib.net
(:require ["@chainsafe/libp2p-gossipsub" :refer (gossipsub)]
["@chainsafe/libp2p-noise" :refer (noise)]
["@chainsafe/libp2p-yamux" :refer (yamux)]
["@libp2p/bootstrap" :refer (bootstrap)]
["@libp2p/circuit-relay-v2" :refer (circuitRelayTransport)]
["@libp2p/crypto/keys" :refer (generateKeyPair)]
["@libp2p/identify" :refer (identify)]
["@libp2p/peer-id" :refer (peerIdFromPrivateKey)]
["@libp2p/ping" :refer (ping)]
["@libp2p/pubsub-peer-discovery" :refer (pubsubPeerDiscovery)]
["@libp2p/webrtc" :refer (webRTC)]
["@libp2p/websockets" :refer (webSockets)]
["@multiformats/multiaddr" :refer (multiaddr)]
["libp2p" :refer (createLibp2p)]
["it-length-prefixed" :as lp]
["it-pushable" :refer (pushable)]
["it-pipe" :refer (pipe)]
[ardegazu.rooms.js :as j]
[ardegazu.rooms.lib.crypto :as rc]
[shadow.cljs.modern :refer (defclass js-await)]))
(def ^:private MSG-PROTOCOL "/sueta/2/msg/1.0")
(def ^:private BEACON-MS 10000) ; presence heartbeat
(def ^:private STALE-MS 30000) ; no decryptable beacon for this long ⇒ gone
(def ^:private STRANGER-DROP-MS 20000) ; connected but never proved membership ⇒ hang up
(def ^:private PING-MS 25000) ; relay liveness probe (mobile TCP half-open)
(def ^:private SWEEP-MS 5000) ; stale/stranger sweep cadence
(def ^:private MAX-FRAME (* 256 1024)) ; matches v1 relay cap; images are chunked below this
;; Throttles. Each one exists to stop a loop, and the comment at its use site
;; says which loop.
(def ^:private DIAL-THROTTLE-MS 15000) ; a dial attempt is still in flight
(def ^:private HELLO-THROTTLE-MS 2000) ; fast-hello beacon vs. beacon ping-pong
(def ^:private BEACON-TO-THROTTLE-MS 5000) ; directed beacon vs. the same
(def ^:private PING-FAIL-LIMIT 2) ; consecutive ping failures before a hang-up
;; Backoff and dial deadlines.
(def ^:private REDIAL-BASE-MS 500)
(def ^:private REDIAL-CAP-MS 15000)
(def ^:private RELAY-DIAL-TIMEOUT-MS 15000)
(def ^:private PEER-DIAL-TIMEOUT-MS 30000)
(def ^:private STREAM-DIAL-TIMEOUT-MS 15000)
(def ^:private PING-TIMEOUT-MS 10000)
(def ^:private te (js/TextEncoder.))
(def ^:private td (js/TextDecoder.))
(defn new-session-key
"Fresh network identity for this session; the relay can't link page loads."
[]
(js-await [private-key (generateKeyPair "Ed25519")]
(j/ordered "privateKey" private-key
"peerId" (.toString (peerIdFromPrivateKey private-key)))))
(defn self-peer-ids
"Every libp2p node this PROCESS runs, by peer id. It is a global so the suite's
kits share ONE registry and sibling nodes never dial each other: a failed
sibling upgrade drives the native WebRTC layer through teardowns that can
abort a Node host. Harmless in a browser, where a page load runs exactly one
Net and the set holds only our own id — which `my-id` already excludes."
[]
(let [g js/globalThis]
(when-not (j/truthy? (unchecked-get g "__ardegazuSelfPeers"))
(unchecked-set g "__ardegazuSelfPeers" (js/Set.)))
(unchecked-get g "__ardegazuSelfPeers")))
(defn- sibling?
"Is `id` another libp2p node inside this same process?"
[id]
(let [^js mine (self-peer-ids)]
(.has mine id)))
;; NetEvents (the object Net.create takes):
;; peerState(peer, state, name) membership-verified peer appeared or its
;; transport state changed
;; peerGone(peer)
;; message(from, payload) decrypted broadcast or directed JSON payload
;; binary(from, data) decrypted directed binary payload
;; peerReady(peer) connected + membership-verified; safe to
;; sendTo (v1 channelOpen)
;; status(up) relay connectivity, for the UI's signaling dot
;; A PeerRec is a mutable JS object (never serialized, so key order is
;; irrelevant — but built with sequential sets by policy, dev/docs/CLJS.md):
;; state PeerConnState
;; name undefined until a beacon names them
;; lastBeacon 0 = never verified
;; firstSeen for the stranger-drop clock
;; lastDialAt dial throttle (upgrade attempts ride the beacon cadence)
;; beaconedAt directed-beacon throttle (pre-upgrade membership proof)
;; ready peerReady already fired
;; out Map of stream protocol -> {push, stream}
(defn- new-peer-rec []
(j/ordered "state" "connecting"
"name" js/undefined
"lastBeacon" 0
"firstSeen" (js/Date.now)
"lastDialAt" 0
"beaconedAt" 0
"ready" false
"out" (js/Map.)))
(defclass Net
(constructor [this node room-crypto events my-name relay-addr]
(unchecked-set this "libp2p" node)
(unchecked-set this "_crypto" room-crypto)
(unchecked-set this "_events" events)
(unchecked-set this "_myName" my-name)
(unchecked-set this "_relayAddr" relay-addr)
;; TS `relayAddr.split("/p2p/").pop() ?? ""` — NULLISH, and .pop on an empty
;; array is undefined, which is the only way the fallback is ever taken
(unchecked-set this "_relayPeer" (j/nn (.pop (.split relay-addr "/p2p/")) ""))
(unchecked-set this "_peers" (js/Map.))
(unchecked-set this "_timers" (array))
(unchecked-set this "_redialTimer" nil)
(unchecked-set this "_redialAttempt" 0)
(unchecked-set this "_windowListeners" (array))
(unchecked-set this "_relayUp" false)
(unchecked-set this "_closed" false)
(unchecked-set this "_lastBeaconAt" 0)))
;; ---- the objects hanging off `self` ----------------------------------------
;;
;; libp2p and the events object are plain JS values reached through
;; `unchecked-get`, and dev/docs/CLJS.md's externs pitfall applies: a `^js` hint
;; placed on a MACRO form (unchecked-get) is lost, so `(.getConnections
;; (unchecked-get self "libp2p"))` cannot be type-inferred. These wrappers bind
;; the object to a `^js`-tagged local first, which is the canon's fix; `ev-call`
;; goes further and reads the callback as a property instead of calling a method,
;; so no interop name is involved at all (peer-kit's idiom).
(defn my-id [self]
(.toString (unchecked-get (unchecked-get self "libp2p") "peerId")))
(defn- ev-call
"Invoke one of the NetEvents callbacks by string key."
[self k & args]
(apply (unchecked-get (unchecked-get self "_events") k) args)
js/undefined)
(defn- crypto-of [self] (unchecked-get self "_crypto"))
(defn- room-topic [self] (unchecked-get (crypto-of self) "roomTopic"))
(defn- peers-of [self] (unchecked-get self "_peers"))
(defn- all-conns [self]
(let [^js node (unchecked-get self "libp2p")]
(.getConnections node)))
(defn- libp2p-dial [self ma opts]
(let [^js node (unchecked-get self "libp2p")]
(.dial node ma opts)))
(defn- libp2p-hang-up [self peer]
(let [^js node (unchecked-get self "libp2p")]
(.hangUp node peer)))
(defn- libp2p-known-peers [self]
(let [^js node (unchecked-get self "libp2p")]
(.getPeers node)))
(defn- libp2p-dial-protocol [self target protocol opts]
(let [^js node (unchecked-get self "libp2p")]
(.dialProtocol node target protocol opts)))
(defn- libp2p-handle [self protocol handler opts]
(let [^js node (unchecked-get self "libp2p")]
(.handle node protocol handler opts)))
(defn- pubsub [self]
(unchecked-get (unchecked-get (unchecked-get self "libp2p") "services") "pubsub"))
(defn- from-peer?
"Is `c` a connection to peer `id`?"
[c id]
(identical? (.toString (unchecked-get c "remotePeer")) id))
(defn- conns-of
"Open connections to `id`."
[self id]
(.filter ^js (all-conns self)
(fn [c] (and (from-peer? c id)
(identical? "open" (unchecked-get c "status"))))))
(defn- direct?
"Does any of these connections ride a /webrtc address?"
[conns]
(.some ^js conns (fn [c] (.includes (.toString (unchecked-get c "remoteAddr")) "/webrtc"))))
(defn- relay-connected?
"Do we hold an OPEN connection to the relay right now?"
[self]
(let [relay-peer (unchecked-get self "_relayPeer")]
(.some ^js (all-conns self)
(fn [c] (and (from-peer? c relay-peer)
(identical? "open" (unchecked-get c "status")))))))
;; ---- decisions, as plain values --------------------------------------------
;;
;; Four rules that were expressions buried mid-effect. Each takes values and
;; returns a value; the shell below is what acts on the answer.
(defn- frame
"One directed frame: a 1-byte kind, then the body."
[kind body]
(let [out (js/Uint8Array. (+ 1 (.-length ^js body)))]
(aset out 0 kind)
(.set out body 1)
out))
(defn- redial-delay
"Jittered exponential backoff, same shape as v1 signaling.ts: 500·2^attempt
capped at 15 s, times a 0.5–1.5 jitter factor."
[attempt jitter]
(* (js/Math.min REDIAL-CAP-MS (* REDIAL-BASE-MS (js/Math.pow 2 attempt)))
(+ 0.5 jitter)))
(defn- conn-state
"The PeerConnState a peer's connections imply. `verified?` only separates the
two zero-connection cases: a peer that once proved membership has
DISCONNECTED; one that never did is still CONNECTING."
[conns verified?]
(cond
(identical? 0 (.-length ^js conns)) (if verified? "disconnected" "connecting")
(direct? conns) "direct"
;; circuit only — privacy unaffected (sealed e2e), just slower
:else "relayed"))
(defn- dial-target
"The multiaddr string to dial `id` at, or nil for \"don't\".
Two-stage dialing (quota + privacy): a DISCOVERED peer is dialed over the bare
relay circuit only — cheap, no ICE, no TURN allocation, and the peer never
sees any of our candidates. Only once it PROVES membership (§4b: something it
sent decrypted) do we dial the WebRTC upgrade, which gathers ICE and — under
`&relay` — consumes a coturn allocation. Strangers from the shared discovery
topic therefore cost nothing and learn nothing; without this, ~7 stranger
dials could exhaust user-quota=16 and starve the call media pcs (seen live as
\"486 TURN allocate error\").
nil when there is nothing to gain: a verified peer already upgraded, an
unverified peer connected at all, or an attempt still in flight."
[relay-addr id verified? conns upgraded? since-dial]
(cond
(if verified? upgraded? (> (.-length ^js conns) 0)) nil
(< since-dial DIAL-THROTTLE-MS) nil
;; member: upgrade to direct
verified? (str relay-addr "/p2p-circuit/webrtc/p2p/" id)
;; unknown: circuit probe only
:else (str relay-addr "/p2p-circuit/p2p/" id)))
(defn- sweep-verdict
"What the stale sweep should do with one peer record. :drop — a verified member
has gone quiet. :forget — a stranger connected (or discovered) but never sent
a decryptable beacon inside the window, so it is not a member of this room and
is disconnected quietly (v1 §3 rule). nil — leave it alone."
[rec now]
(if (j/truthy? (unchecked-get rec "lastBeacon"))
(when (> (- now (unchecked-get rec "lastBeacon")) STALE-MS) :drop)
(when (> (- now (unchecked-get rec "firstSeen")) STRANGER-DROP-MS) :forget)))
;; ---- relay connection management -------------------------------------------
(defn- set-relay-up [self up]
(when-not (identical? (unchecked-get self "_relayUp") up)
(unchecked-set self "_relayUp" up)
(ev-call self "status" up))
js/undefined)
;; The one real cycle in this namespace: a failed relay dial schedules a retry,
;; and the retry dials again.
(declare ensure-relay)
(defn- schedule-redial
"Queue one relay retry, unless one is already queued or we are closed. A timer
id of 0 is FALSY in JS and means \"nothing scheduled\" — so this guard is JS
truthiness, not Clojure's, and the harness returns exactly that 0."
[self]
(when-not (or (j/truthy? (unchecked-get self "_redialTimer"))
(j/truthy? (unchecked-get self "_closed")))
(let [attempt (unchecked-get self "_redialAttempt")
ms (redial-delay attempt (js/Math.random))]
(unchecked-set self "_redialAttempt" (inc attempt))
(unchecked-set self "_redialTimer"
(js/setTimeout (fn []
(unchecked-set self "_redialTimer" nil)
(ensure-relay self))
ms))))
js/undefined)
(defn- ensure-relay
"Make sure we hold a relay connection, dialing — and scheduling a retry — if
not. A successful dial resets the backoff ladder."
[self]
(when-not (j/truthy? (unchecked-get self "_closed"))
(if (relay-connected? self)
(set-relay-up self true)
(do
(set-relay-up self false)
(-> (libp2p-dial self (multiaddr (unchecked-get self "_relayAddr"))
(j/ordered "signal" (js/AbortSignal.timeout RELAY-DIAL-TIMEOUT-MS)))
(.then (fn [_]
(unchecked-set self "_redialAttempt" 0)
(set-relay-up self true)))
(.catch (fn [_] (schedule-redial self)))))))
js/undefined)
(defn- relay-peer-id
"The relay's libp2p PeerId object, from any connection to it; nil when there is
none. Deliberately NOT status-filtered — ping wants whatever handle exists."
[self]
(let [relay-peer (unchecked-get self "_relayPeer")
conn (.find ^js (all-conns self) (fn [c] (from-peer? c relay-peer)))]
(when (j/truthy? conn) (unchecked-get conn "remotePeer"))))
;; ---- peer lifecycle --------------------------------------------------------
(defn- rec-for
"This peer's record, created on first sight — which is what starts its
stranger-drop clock."
[self id]
(let [^js peers (peers-of self)]
(or (.get peers id)
(let [fresh (new-peer-rec)]
(.set peers id fresh)
fresh))))
;; `_self` — receiver kept for the family's shape; nothing is read off it here.
(defn- end-channels [_self rec]
(.forEach ^js (unchecked-get rec "out") (fn [chan _k] (.end (unchecked-get chan "push"))))
(.clear ^js (unchecked-get rec "out"))
js/undefined)
(defn- hang-up-str [self id]
(let [conn (.find ^js (all-conns self) (fn [c] (from-peer? c id)))]
(when (j/truthy? conn)
(.catch (libp2p-hang-up self (unchecked-get conn "remotePeer"))
(fn [_] nil))))
js/undefined)
(defn- drop-peer [self id rec]
(end-channels self rec)
(.delete ^js (peers-of self) id)
(hang-up-str self id)
(when (j/truthy? (unchecked-get rec "lastBeacon"))
(ev-call self "peerGone" id))
js/undefined)
(defn- refresh-state
"Recompute a peer's transport state and, when it changed, tell the app. Losing
the last connection also tears down the outbound channels and clears `ready`."
[self id rec emit?]
(let [conns (conns-of self id)
verified? (j/truthy? (unchecked-get rec "lastBeacon"))
state (conn-state conns verified?)]
(when (identical? 0 (.-length ^js conns))
(unchecked-set rec "ready" false)
(end-channels self rec))
(when-not (identical? state (unchecked-get rec "state"))
(unchecked-set rec "state" state)
(when (and emit? verified?)
(ev-call self "peerState" id state (unchecked-get rec "name")))))
js/undefined)
(defn- maybe-ready
"peerReady = live connection + proven membership, fired once per session."
[self id rec]
(let [has-conn (> (.-length ^js (conns-of self id)) 0)]
(when (and (not (j/truthy? (unchecked-get rec "ready")))
has-conn
(j/truthy? (unchecked-get rec "lastBeacon")))
(unchecked-set rec "ready" true)
(ev-call self "peerReady" id)))
js/undefined)
(defn open-peers [self]
(let [out (array)]
(.forEach ^js (peers-of self)
(fn [rec id] (when (j/truthy? (unchecked-get rec "ready")) (.push out id))))
out))
;; ---- directed streams (generic per-protocol frame channels) ----------------
(defn handle-frames
"Register a length-prefixed frame handler for a stream protocol. Frames are raw
bytes — the caller does its own sealing/opening."
[self protocol on-frame]
(libp2p-handle
self
protocol
(fn [conn-info]
(let [stream (unchecked-get conn-info "stream")
from (.toString (unchecked-get (unchecked-get conn-info "connection") "remotePeer"))]
;; stream reset — sender reopens on next send
(.catch (pipe stream
(fn [s] (lp/decode s (j/ordered "maxDataLength" MAX-FRAME)))
(fn [frames]
(j/for-each! frames (fn [f] (on-frame from (.subarray ^js f)) nil))))
(fn [_] nil))
js/undefined))
;; sealed frames must work over bare relay circuits (pre-upgrade members, the
;; "relayed" fallback) — the crypto layer, not the transport, gates trust
(j/ordered "runOnLimitedConnection" true)))
(defn- live-channel
"The peer's existing outbound stream for `protocol`, if it is still open."
[rec protocol]
(let [chan (.get ^js (unchecked-get rec "out") protocol)]
(when (and (j/truthy? chan)
(identical? "open" (unchecked-get (unchecked-get chan "stream") "status")))
chan)))
(defn- stream-target
"Where to open a stream to `peer`: a PeerId libp2p already knows, else an
explicit relay-circuit multiaddr."
[self peer]
(let [pid (.find ^js (libp2p-known-peers self) (fn [x] (identical? (.toString x) peer)))]
(if (some? pid)
pid
(multiaddr (str (unchecked-get self "_relayAddr") "/p2p-circuit/webrtc/p2p/" peer)))))
(defn- open-channel!
"Open a fresh outbound stream for (peer, protocol), register it on the peer
record and push the first frame down it. Resolves true."
[self rec protocol target frame-bytes]
(js-await [stream (libp2p-dial-protocol
self target protocol
(j/ordered "signal" (js/AbortSignal.timeout STREAM-DIAL-TIMEOUT-MS)
"runOnLimitedConnection" true))]
(let [push (pushable)
^js out (unchecked-get rec "out")
chan (j/ordered "push" push "stream" stream)]
(.set out protocol chan)
;; the stream died: forget it so the next send reopens — unless another
;; channel has already replaced this one
(.catch (pipe push (fn [s] (lp/encode s)) stream)
(fn [_]
(when (identical? (.get out protocol) chan)
(.delete out protocol))
nil))
(.push ^js push frame-bytes)
true)))
(defn send-frame
"Send one frame on a per-(peer, protocol) persistent outbound stream; resolves
false when the peer is unreachable."
[self protocol peer frame-bytes]
;; transport-level: no membership gate here — replies to decrypted directed
;; messages may target peers whose beacon hasn't arrived yet
(let [rec (rec-for self peer)
chan (live-channel rec protocol)]
(if (some? chan)
(do (.push ^js (unchecked-get chan "push") frame-bytes)
(js/Promise.resolve true))
;; the target is resolved EAGERLY, as the original did: a malformed relay
;; multiaddr is a configuration error and must not be swallowed as "peer
;; unreachable"
(let [target (stream-target self peer)]
(j/attempt (fn [] (open-channel! self rec protocol target frame-bytes)) false)))))
;; ---- sealed senders --------------------------------------------------------
(defn broadcast
"Seal once, deliver to every connected member via the room topic."
[self payload]
(js-await [env (rc/seal-msg (crypto-of self) payload)]
(let [bytes (.encode te (js/JSON.stringify env))
^js ps (pubsub self)]
;; alone in the room is fine
(.catch (js/Promise.resolve (.publish ps (room-topic self) bytes))
(fn [_] nil)))))
(defn send-to
"Directed JSON payload; resolves false if the peer is unreachable."
[self peer payload]
(js-await [env (rc/seal-msg (crypto-of self) payload)]
(send-frame self MSG-PROTOCOL peer
(frame 0x00 (.encode te (js/JSON.stringify env))))))
(defn send-binary-to [self peer data]
(js-await [sealed (rc/seal-msg-binary (crypto-of self) data)]
(js-await [_ (send-frame self MSG-PROTOCOL peer (frame 0x01 (js/Uint8Array. sealed)))]
js/undefined)))
(defn broadcast-binary
"Binary fan-out to every ready peer. Seals once."
[self data]
(js-await [sealed-buf (rc/seal-msg-binary (crypto-of self) data)]
(let [f (frame 0x01 (js/Uint8Array. sealed-buf))]
(js-await [_ (js/Promise.all
(.map ^js (open-peers self)
(fn [id] (send-frame self MSG-PROTOCOL id f))))]
js/undefined))))
;; ---- dialing, presence, membership -----------------------------------------
(defn- dial-peer
"Dial `id`, if `dial-target` says there is anything to gain."
[self id]
(if (j/truthy? (unchecked-get self "_closed"))
(js/Promise.resolve nil)
(let [id-str (.toString id)
rec (rec-for self id-str) ; starts the stranger-drop clock
now (js/Date.now)
conns (conns-of self id-str)
target (dial-target (unchecked-get self "_relayAddr") id-str
(> (unchecked-get rec "lastBeacon") 0)
conns (direct? conns)
(- now (unchecked-get rec "lastDialAt")))]
(if (nil? target)
(js/Promise.resolve nil)
(do
(unchecked-set rec "lastDialAt" now)
(-> (libp2p-dial self (multiaddr target)
(j/ordered "signal" (js/AbortSignal.timeout PEER-DIAL-TIMEOUT-MS)))
;; the other side may dial us instead
(.catch (fn [_] nil))
(.then (fn [_] nil))))))))
(defn- presence-frame
"The presence payload, in wire key order — it is sealed and JSON.stringify'd,
and the topic beacon and the directed beacon must build the identical shape."
[self op]
(j/ordered "kind" "presence"
"op" op
"name" ((unchecked-get self "_myName"))
"ts" (js/Date.now)))
(defn- beacon
"Announce ourselves on the room topic. `op` is \"beacon\" or \"bye\"."
[self op]
(if (j/truthy? (unchecked-get self "_closed"))
(js/Promise.resolve nil)
(do
(when (identical? "beacon" op) (unchecked-set self "_lastBeaconAt" (js/Date.now)))
(.catch (broadcast self (presence-frame self op)) (fn [_] nil)))))
(defn- beacon-to
"Directed sealed beacon: the pre-upgrade membership proof. Gossipsub only runs
on direct connections, so a peer we just met over a bare circuit would never
see our topic beacons — this hands them one over the msg stream (which does
run on limited connections). Decrypting it proves OUR membership; their
directed reply proves theirs; then both sides upgrade."
[self peer]
(let [rec (rec-for self peer)]
(when-not (< (- (js/Date.now) (unchecked-get rec "beaconedAt")) BEACON-TO-THROTTLE-MS)
(unchecked-set rec "beaconedAt" (js/Date.now))
(send-to self peer (presence-frame self "beacon"))))
js/undefined)
(defn- verified
"Any payload that decrypts under the room keys IS proof of membership (v1 §3:
decryptability is the trust boundary). Called for every decrypted inbound
payload so directed messages from a peer whose beacon hasn't arrived yet still
verify it — e.g. a hist-req right after joining."
[self from]
(let [rec (rec-for self from)
first? (identical? 0 (unchecked-get rec "lastBeacon"))]
(unchecked-set rec "lastBeacon" (js/Date.now))
(refresh-state self from rec (not first?))
(when first?
(ev-call self "peerState" from
(unchecked-get rec "state") (unchecked-get rec "name"))
(maybe-ready self from rec)
;; fast hello: answer with our own beacon so the new member verifies us
;; without waiting out the 10 s cadence (throttled against beacon loops);
;; the directed copy reaches them even while we're still circuit-only
(when (> (- (js/Date.now) (unchecked-get self "_lastBeaconAt")) HELLO-THROTTLE-MS)
(beacon self "beacon"))
(beacon-to self from)
;; membership proven — NOW spend the WebRTC upgrade on this peer
(unchecked-set rec "lastDialAt" 0)
(dial-peer self from))
rec))
(defn- on-presence
"Presence handling, shared by the room topic and the directed msg stream."
[self from p]
(if (identical? "bye" (unchecked-get p "op"))
(let [known (.get ^js (peers-of self) from)]
(when (and (some? known) (j/truthy? (unchecked-get known "lastBeacon")))
(drop-peer self from known))
js/undefined)
(let [prev (.get ^js (peers-of self) from)
was-known (> (if (some? prev) (j/nn (unchecked-get prev "lastBeacon") 0) 0) 0)
prev-name (when (some? prev) (unchecked-get prev "name"))
pname (unchecked-get p "name")]
(when (identical? "string" (js* "typeof ~{}" pname))
(unchecked-set (rec-for self from) "name" pname))
(let [rec (verified self from)]
(when (and was-known (not (identical? (unchecked-get rec "name") prev-name)))
(ev-call self "peerState" from
(unchecked-get rec "state") (unchecked-get rec "name")))
;; keep verified members converging on a direct connection: dial-peer
;; no-ops when already upgraded and self-throttles between attempts
(dial-peer self from)
js/undefined))))
(defn- on-sealed-json
"The inbound path for a sealed JSON envelope, whichever transport carried it —
the room topic and the directed 0x00 frame ran byte-identical code. Both the
parse and the decrypt are drop-on-failure: these are network bytes.
Directed presence is the pre-upgrade membership handshake, so it is handled
like topic presence and never surfaced to the app."
[self from bytes]
(let [env (j/parse-json (.decode td bytes))]
(if (undefined? env)
(js/Promise.resolve nil)
(js-await [payload (rc/open-msg (crypto-of self) from env)]
(if (nil? payload)
nil ; not a member (or tampered) — drop silently
(if (identical? "presence" (unchecked-get payload "kind"))
(do (on-presence self from payload) nil)
(do (verified self from)
(ev-call self "message" from payload)
nil)))))))
(defn- on-topic-message
"Room-topic delivery. Gossipsub only binds to direct connections, so this is
the post-upgrade path; beacon-to carries the same envelopes before that."
[self from data]
(on-sealed-json self from data))
(defn- on-msg-frame
"One directed frame: 0x00 is a sealed JSON envelope, 0x01 a sealed binary
payload, anything else — including a frame too short to have a body — a drop."
[self from buf]
(if (< (.-length ^js buf) 2)
(js/Promise.resolve nil)
(let [kind (aget buf 0)
body (.subarray ^js buf 1)]
(cond
(identical? 0x00 kind) (on-sealed-json self from body)
(identical? 0x01 kind)
(let [ab (js/ArrayBuffer. (.-length body))]
(.set (js/Uint8Array. ab) body)
(js-await [data (rc/open-msg-binary (crypto-of self) from ab)]
(if (nil? data)
nil
(do (verified self from)
(ev-call self "binary" from data)
nil))))
:else (js/Promise.resolve nil)))))
(defn- on-conn-change [self conn]
(let [id (.toString (unchecked-get conn "remotePeer"))]
(if (identical? id (unchecked-get self "_relayPeer"))
(let [still-up (relay-connected? self)]
(set-relay-up self still-up)
(when-not still-up (schedule-redial self)))
(let [rec (.get ^js (peers-of self) id)]
(when (some? rec)
(refresh-state self id rec true)
(maybe-ready self id rec)
;; fresh contact over a bare circuit: hand them our membership proof
;; (gossipsub can't — it only runs on direct connections)
(when (and (identical? "open" (unchecked-get conn "status"))
(not (j/truthy? (unchecked-get rec "ready"))))
(beacon-to self id))))))
js/undefined)
(defn- sweep [self]
(let [now (js/Date.now)]
(.forEach ^js (peers-of self)
(fn [rec id]
(case (sweep-verdict rec now)
:drop (drop-peer self id rec)
:forget (do (.delete ^js (peers-of self) id)
(hang-up-str self id))
nil))))
js/undefined)
;; ---- app-facing helpers ----------------------------------------------------
(defn debug-state [self]
(let [out (array)]
(.forEach ^js (peers-of self)
(fn [rec id]
(.push out (j/ordered
"id" id
"state" (unchecked-get rec "state")
"name" (unchecked-get rec "name")
"ready" (unchecked-get rec "ready")
"verified" (> (unchecked-get rec "lastBeacon") 0)
"conns" (.map ^js (.filter ^js (all-conns self)
(fn [c] (from-peer? c id)))
(fn [c] (.toString (unchecked-get c "remoteAddr"))))))))
out))
(defn close [self]
(unchecked-set self "_closed" true)
(js-await [_ (.catch (beacon self "bye") (fn [_] nil))]
(do
(doseq [pair (array-seq (unchecked-get self "_windowListeners"))]
(.removeEventListener js/window (aget pair 0) (aget pair 1)))
(unchecked-set self "_windowListeners" (array))
(doseq [t (array-seq (unchecked-get self "_timers"))] (js/clearInterval t))
(when (j/truthy? (unchecked-get self "_redialTimer"))
(js/clearTimeout (unchecked-get self "_redialTimer")))
(.forEach ^js (peers-of self) (fn [rec _id] (end-channels self rec)))
(js-await [_ (.catch (js/Promise.resolve (.stop ^js (unchecked-get self "libp2p")))
(fn [_] nil))]
js/undefined))))
;; ---- start: one installer per subsystem ------------------------------------
(defn- install-discovery!
"Dial peers the relay's shared discovery topic surfaces — but never the relay
itself, never ourselves, and never a SIBLING node in this same process."
[self]
(let [^js node (unchecked-get self "libp2p")]
(.addEventListener node "peer:discovery"
(fn [e]
(let [detail (unchecked-get e "detail")
id (.toString (unchecked-get detail "id"))]
(when-not (or (identical? id (unchecked-get self "_relayPeer"))
(identical? id (my-id self))
(sibling? id))
(dial-peer self (unchecked-get detail "id")))
js/undefined))))
js/undefined)
(defn- install-conn-listeners! [self]
(let [^js node (unchecked-get self "libp2p")
on-change (fn [e] (on-conn-change self (unchecked-get e "detail")))]
(.addEventListener node "connection:open" on-change)
(.addEventListener node "connection:close" on-change))
js/undefined)
(defn- install-topic!
"Subscribe to the room topic and route its messages. Another topic's traffic
and our own echo are dropped before decryption is even attempted."
[self]
(let [^js ps (pubsub self)
topic (room-topic self)]
(.subscribe ps topic)
(.addEventListener ps "message"
(fn [e]
(let [detail (unchecked-get e "detail")]
(when (identical? (unchecked-get detail "topic") topic)
(let [f (unchecked-get detail "from")
from (when (some? f) (.toString f))]
(when (and (j/truthy? from)
(not (identical? from (my-id self))))
(on-topic-message self from (unchecked-get detail "data")))))
js/undefined))))
js/undefined)
(defn- ping-relay!
"One relay liveness probe. PING-FAIL-LIMIT consecutive failures hang the relay
connection up, which drives on-conn-change into a redial — mobile TCP goes
half-open silently and a ping that never returns is the only signal.
`fails` is a cell owned by the timer that calls this, not a field on `self`:
it is private to that one interval and dies with it."
[self fails]
(when (and (j/truthy? (unchecked-get self "_relayUp"))
(not (j/truthy? (unchecked-get self "_closed"))))
(let [relay (relay-peer-id self)]
(when (j/truthy? relay)
(-> (let [^js svc (unchecked-get (unchecked-get (unchecked-get self "libp2p") "services") "ping")]
(.ping svc relay (j/ordered "signal" (js/AbortSignal.timeout PING-TIMEOUT-MS))))
(.then (fn [_] (vreset! fails 0) nil))
(.catch (fn [_]
(vswap! fails inc)
(when (>= @fails PING-FAIL-LIMIT)
(vreset! fails 0)
(.catch (libp2p-hang-up self relay) (fn [_] nil)))
nil))))))
js/undefined)
(defn- install-timers!
"Presence heartbeat, stale sweep and relay liveness ping. Every id is kept so
close() can clear it."
[self]
(let [^js timers (unchecked-get self "_timers")
ping-fails (volatile! 0)]
(.push timers (js/setInterval (fn [] (beacon self "beacon")) BEACON-MS))
(.push timers (js/setInterval (fn [] (sweep self)) SWEEP-MS))
(.push timers (js/setInterval (fn [] (ping-relay! self ping-fails)) PING-MS)))
js/undefined)
(defn- install-wake-listeners!
"Mobile wake: check the relay connection immediately and re-announce. Handlers
are kept so close() can remove them — a consumer that creates and closes Nets
without a page reload must not leak them (each closure retains the whole
node)."
[self]
(let [kick (fn []
(when-not (j/truthy? (unchecked-get self "_closed"))
(unchecked-set self "_redialAttempt" 0)
(ensure-relay self)
(beacon self "beacon"))
js/undefined)
listeners #js [#js ["visibilitychange"
(fn [] (when (identical? "visible" (.-visibilityState js/document)) (kick)))]
#js ["pageshow" kick]
#js ["online" kick]
#js ["pagehide" (fn [] (beacon self "bye") js/undefined)]]]
(unchecked-set self "_windowListeners" listeners)
(doseq [pair (array-seq listeners)]
(.addEventListener js/window (aget pair 0) (aget pair 1))))
js/undefined)
(defn- start [self]
(js-await [_ (handle-frames self MSG-PROTOCOL
(fn [from frm] (on-msg-frame self from frm) nil))]
(do
(install-discovery! self)
(install-conn-listeners! self)
(install-topic! self)
(install-timers! self)
(install-wake-listeners! self)
(ensure-relay self)
(beacon self "beacon")
js/undefined)))
(defn create [opts]
(js-await [node (createLibp2p
(j/ordered
"privateKey" (unchecked-get opts "privateKey")
"addresses" (j/ordered "listen" #js ["/p2p-circuit" "/webrtc"])
;; `wsOrigin` — a Node host must send an Origin header to
;; reach the origin-gated managed relay. Absent (the browser,
;; which the user agent already originates) this is
;; `webSockets()` with no options, byte for byte.
"transports" #js [(let [origin (unchecked-get opts "wsOrigin")]
(webSockets (if (j/truthy? origin)
(j/ordered "websocket" (j/ordered "origin" origin))
js/undefined)))
(webRTC (j/ordered "rtcConfiguration"
(fn [] (.get (unchecked-get opts "ice")))))
(circuitRelayTransport)]
"connectionEncrypters" #js [(noise)]
"streamMuxers" #js [(yamux)]
"peerDiscovery" #js [(bootstrap (j/ordered "list" #js [(unchecked-get opts "relayMultiaddr")]))
(pubsubPeerDiscovery
(j/ordered "interval" 5000
"topics" #js [(unchecked-get opts "discoveryTopic")]))]
"services" (j/ordered
"identify" (identify)
"ping" (ping)
;; gossipsub deliberately does NOT run on limited
;; (circuit) connections: its streams must only
;; ever bind to direct webrtc connections, so an
;; established mesh keeps flowing when the relay
;; dies. Pre-upgrade membership proof rides
;; DIRECTED beacons on the msg stream instead
;; (handle-frames/send-frame run on limited
;; connections).
"pubsub" (gossipsub (j/ordered "allowPublishToZeroTopicPeers" true
"emitSelf" false)))
"connectionManager" (j/ordered "maxConnections" 32)))]
(let [net (Net. node (unchecked-get opts "crypto") (unchecked-get opts "events")
(unchecked-get opts "myName") (unchecked-get opts "relayMultiaddr"))]
;; never dial our own sibling nodes
(.add ^js (self-peer-ids) (.toString (unchecked-get node "peerId")))
(js-await [_ (start net)]
net))))
;; ---- class surface ---------------------------------------------------------
(unchecked-set Net "create" (fn [opts] (create opts)))
(let [proto (.-prototype Net)]
(js/Object.defineProperty
proto "myId"
(j/ordered "get" (fn [] (this-as self (my-id self))) "configurable" true))
(js/Object.defineProperty
proto "relayUp"
(j/ordered "get" (fn [] (this-as self (unchecked-get self "_relayUp"))) "configurable" true))
(unchecked-set proto "handleFrames"
(fn [protocol on-frame] (this-as self (handle-frames self protocol on-frame))))
(unchecked-set proto "sendFrame"
(fn [protocol peer f] (this-as self (send-frame self protocol peer f))))
(unchecked-set proto "sendTo" (fn [peer payload] (this-as self (send-to self peer payload))))
(unchecked-set proto "sendBinaryTo" (fn [peer data] (this-as self (send-binary-to self peer data))))
(unchecked-set proto "broadcast" (fn [payload] (this-as self (broadcast self payload))))
(unchecked-set proto "broadcastBinary" (fn [data] (this-as self (broadcast-binary self data))))
(unchecked-set proto "openPeers" (fn [] (this-as self (open-peers self))))
(unchecked-set proto "debugState" (fn [] (this-as self (debug-state self))))
(unchecked-set proto "close" (fn [] (this-as self (close self))))
;; the "private" surface the tests drive (peer-kit's `_name` idiom)
(unchecked-set proto "_dialPeer" (fn [id] (this-as self (dial-peer self id))))
(unchecked-set proto "_recFor" (fn [id] (this-as self (rec-for self id))))
(unchecked-set proto "_onConnChange" (fn [conn] (this-as self (on-conn-change self conn))))
(unchecked-set proto "_refreshState" (fn [id rec emit?] (this-as self (refresh-state self id rec emit?))))
(unchecked-set proto "_maybeReady" (fn [id rec] (this-as self (maybe-ready self id rec))))
(unchecked-set proto "_sweep" (fn [] (this-as self (sweep self))))
(unchecked-set proto "_endChannels" (fn [rec] (this-as self (end-channels self rec))))
(unchecked-set proto "_dropPeer" (fn [id rec] (this-as self (drop-peer self id rec))))
(unchecked-set proto "_hangUpStr" (fn [id] (this-as self (hang-up-str self id))))
(unchecked-set proto "_beacon" (fn [op] (this-as self (beacon self op))))
(unchecked-set proto "_verified" (fn [from] (this-as self (verified self from))))
(unchecked-set proto "_onTopicMessage" (fn [from data] (this-as self (on-topic-message self from data))))
(unchecked-set proto "_onPresence" (fn [from p] (this-as self (on-presence self from p))))
(unchecked-set proto "_beaconTo" (fn [peer] (this-as self (beacon-to self peer))))
(unchecked-set proto "_onMsgFrame" (fn [from buf] (this-as self (on-msg-frame self from buf))))
(unchecked-set proto "_frame" (fn [kind body] (frame kind body)))
(unchecked-set proto "_relayPeerId" (fn [] (this-as self (relay-peer-id self))))
(unchecked-set proto "_ensureRelay" (fn [] (this-as self (ensure-relay self))))
(unchecked-set proto "_scheduleRedial" (fn [] (this-as self (schedule-redial self))))
(unchecked-set proto "_setRelayUp" (fn [up] (this-as self (set-relay-up self up))))
(unchecked-set proto "_pubsub" (fn [] (this-as self (pubsub self)))))
|