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 | ;; ported-from: src/core/peers.ts @ v1.3.0
;;
;; Peer lifecycle: room-scoped protocol streams, the sealed-hello membership
;; handshake, reconnection, and the verified-member set the app sees.
;;
;; Node port of the suite's shared game net layer (multi-play client/src/net/
;; peers.ts). Two deliberate differences from the browser original:
;; - the protocol prefix is a constructor argument (the browser copies differ
;; by exactly that one line per game);
;; - timers are plain setTimeout (no window), typed for Node and browsers.
;;
;; Rules of the road:
;; - discovery is app-wide; peers from other rooms fail protocol negotiation
;; (different protocol id) and are never retried.
;; - the lexicographically lower peer id dials; the higher side dials anyway
;; after a fallback delay so one-sided discovery still connects.
;; - a stream is trusted only after the remote's hello decrypts under the room
;; secret; an invalid hello blacklists the peer for the session.
;; - a member's stream dying starts quiet re-dials; the app only hears
;; `peerGone` when those are exhausted (signaling loss ≠ peer loss).
(ns ardegazu.peer.core.peers
(:require ["@libp2p/utils" :refer (LengthPrefixedDecoder)]
[ardegazu.peer.obj :as obj]
[ardegazu.peer.core.room-crypto :as rc]
[ardegazu.peer.node.env :as env]
[shadow.cljs.modern :refer (defclass js-await)]))
(def ^:private HELLO-TIMEOUT-MS 10000)
(def ^:private FALLBACK-DIAL-MS 8000)
(def ^:private REDIAL-DELAYS-MS #js [2000 5000 8000])
(def ^:private WAKE-WATCHDOG-MS 12000)
(def ^:private MAX-FRAME-BYTES (* 64 1024))
(def ^:private te (js/TextEncoder.))
(def ^:private td (js/TextDecoder.))
(defn lp-encode
"Varint length prefix, mirroring @libp2p/utils' LengthPrefixedDecoder default."
[bytes]
(let [prefix (array)]
(loop [n (.-length ^js bytes)]
(if (>= n 0x80)
(do (.push prefix (bit-or (bit-and n 0x7f) 0x80))
(recur (unsigned-bit-shift-right n 7)))
(.push prefix n)))
(let [out (js/Uint8Array. (+ (.-length prefix) (.-length ^js bytes)))]
(.set out prefix 0)
(.set out bytes (.-length prefix))
out)))
;; A PeerState is a mutable JS object with these slots (never serialized, so
;; key order is irrelevant here):
;; id peerId stream verified frameChain ident member dialing redialAttempt
;; helloTimer fallbackTimer redialTimer
;; `frameChain` serializes frame handling per peer: openHello is async, and a
;; frame sent right after the hello (a "hi" from a fast local peer) must not
;; race the pending verification — unserialized it would be mistaken for the
;; hello and get the peer permanently rejected. (Latent in the browsers' net
;; too, where real network latency hides it; a candidate upstream fix.)
(defn- new-peer-state [id peer-id]
;; 12 slots — built with sequential sets (ardegazu.peer.obj); internal state,
;; but the eight-pair js-obj rule applies to every wide object by policy
(obj/ordered "id" id
"peerId" peer-id
"stream" nil
"verified" false
"frameChain" (js/Promise.resolve nil)
"ident" nil
"member" false
"dialing" false
"redialAttempt" 0
"helloTimer" 0
"fallbackTimer" 0
"redialTimer" 0))
(defclass PeerManager
(constructor [this node room-crypto relay-id protocol-prefix ev hello-id]
(unchecked-set this "_node" node)
(unchecked-set this "_rc" room-crypto)
(unchecked-set this "_relayId" relay-id)
(unchecked-set this "_ev" ev)
(unchecked-set this "_helloId" (if (identical? hello-id js/undefined) nil hello-id))
(unchecked-set this "_myId" (.toString (.-peerId ^js node)))
(unchecked-set this "_peers" (js/Map.))
(unchecked-set this "_rejected" (js/Set.)) ; bad hello — not a member, never again
(unchecked-set this "_wrongRoom" (js/Set.)) ; doesn't speak our protocol — different room
(unchecked-set this "protocol"
(str protocol-prefix (.slice (unchecked-get room-crypto "roomId") 0 16)))))
(declare pm-dial pm-drop pm-adopt pm-on-frame pm-schedule-redial pm-emit-members
pm-on-stream-closed)
(defn- ev-call
"Invoke an optional callback on the events object (`ev.peerIdentity?.(…)`)."
[self k & args]
(let [f (unchecked-get (unchecked-get self "_ev") k)]
(when ^boolean (js* "!!(~{})" f) (apply f args))
js/undefined))
(defn members [self]
(let [out (array)]
(.forEach ^js (unchecked-get self "_peers")
(fn [p _k] (when (unchecked-get p "member") (.push out (unchecked-get p "id")))))
out))
(defn open-peers [self]
(let [out (array)]
(.forEach ^js (unchecked-get self "_peers")
(fn [p _k]
(when (and (unchecked-get p "member") (unchecked-get p "verified"))
(.push out (unchecked-get p "id")))))
out))
(defn open-count [self]
(.-length (open-peers self)))
(defn identity-of
"Verified suite identity of a member, if it announced one."
[self id]
(let [p (.get ^js (unchecked-get self "_peers") id)]
(if (some? p) (or (unchecked-get p "ident") nil) nil)))
(defn- pm-write [self stream frame]
(try
(.send ^js stream (lp-encode (.encode te (js/JSON.stringify frame))))
(catch :default _ nil)) ; dying stream — the close handler deals with it
js/undefined)
(defn send-to [self id frame]
(let [p (.get ^js (unchecked-get self "_peers") id)]
(when (and (some? p)
^boolean (js* "!!(~{})" (unchecked-get p "stream"))
(unchecked-get p "verified"))
(pm-write self (unchecked-get p "stream") frame)))
js/undefined)
(defn broadcast [self frame]
(let [bytes (lp-encode (.encode te (js/JSON.stringify frame)))]
(.forEach ^js (unchecked-get self "_peers")
(fn [p _k]
(when (and ^boolean (js* "!!(~{})" (unchecked-get p "stream"))
(unchecked-get p "verified"))
(try
(.send ^js (unchecked-get p "stream") bytes)
(catch :default _ nil)))))) ; dying stream — close handles it
js/undefined)
(defn wake
"Post-suspend kick: re-dial members whose streams died while we slept."
[self]
(.forEach ^js (unchecked-get self "_peers")
(fn [p _k]
(when (and (unchecked-get p "member")
(not ^boolean (js* "!!(~{})" (unchecked-get p "stream")))
(not (unchecked-get p "dialing"))
(identical? 0 (unchecked-get p "redialTimer")))
(pm-dial self p))))
(js/setTimeout
(fn []
(let [snapshot (array)]
(.forEach ^js (unchecked-get self "_peers") (fn [p _k] (.push snapshot p)))
(.forEach snapshot
(fn [p]
(when (and (unchecked-get p "member")
(not (unchecked-get p "verified"))
(not (unchecked-get p "dialing")))
(pm-drop self p))))))
WAKE-WATCHDOG-MS)
js/undefined)
(defn- pm-state [self peer-id]
(let [id (.toString ^js peer-id)
peers (unchecked-get self "_peers")]
(or (.get ^js peers id)
(let [p (new-peer-state id peer-id)]
(.set ^js peers id p)
p))))
(defn- pm-skip? [self id]
(or (identical? id (unchecked-get self "_myId"))
(identical? id (unchecked-get self "_relayId"))
(.has (env/self-peer-ids) id) ; a sibling node of this process
(.has ^js (unchecked-get self "_rejected") id)
(.has ^js (unchecked-get self "_wrongRoom") id)))
(defn- pm-on-discovered [self peer-id]
(let [id (.toString ^js peer-id)]
(when-not (pm-skip? self id)
(let [p (pm-state self peer-id)]
(when-not (or ^boolean (js* "!!(~{})" (unchecked-get p "stream"))
(unchecked-get p "dialing")
(not (identical? 0 (unchecked-get p "redialTimer"))))
(if (< (unchecked-get self "_myId") id)
(pm-dial self p)
(when (identical? 0 (unchecked-get p "fallbackTimer"))
;; the lower side dials; give it a head start, then dial anyway
(unchecked-set
p "fallbackTimer"
(js/setTimeout
(fn []
(unchecked-set p "fallbackTimer" 0)
(when-not (or ^boolean (js* "!!(~{})" (unchecked-get p "stream"))
(unchecked-get p "dialing"))
(pm-dial self p)))
FALLBACK-DIAL-MS))))))))
js/undefined)
(defn- pm-dial [self p]
(if (or ^boolean (js* "!!(~{})" (unchecked-get p "stream")) (unchecked-get p "dialing"))
(js/Promise.resolve nil)
(do
(unchecked-set p "dialing" true)
(-> (js/Promise.resolve nil)
(.then
(fn [_]
;; bind to a ^js local first: a ^js hint placed on an
;; `unchecked-get` macro form is LOST (dev/docs/CLJS.md)
(let [^js node (unchecked-get self "_node")]
(js-await [stream (.dialProtocol node
(unchecked-get p "peerId")
(unchecked-get self "protocol")
(js-obj "runOnLimitedConnection" true
"signal" (js/AbortSignal.timeout 15000)))]
(pm-adopt self p stream)))))
(.catch
(fn [err]
(js/console.debug "[net] dial" (.slice (unchecked-get p "id") -6) "failed:"
(when (some? err) (unchecked-get err "name"))
(when (some? err) (unchecked-get err "message")))
(cond
(identical? "UnsupportedProtocolError" (when (some? err) (unchecked-get err "name")))
;; different room on the shared app — not an error
(do (.add ^js (unchecked-get self "_wrongRoom") (unchecked-get p "id"))
(.delete ^js (unchecked-get self "_peers") (unchecked-get p "id")))
(unchecked-get p "member")
(pm-schedule-redial self p)
;; non-member dial failures: the next discovery broadcast retries
:else nil)))
(.finally (fn [] (unchecked-set p "dialing" false)))))))
(defn- pm-on-inbound [self stream conn]
(let [id (.toString (.-remotePeer ^js conn))]
(if (pm-skip? self id)
(do
(when ^boolean (js* "!!(~{})" (unchecked-get (unchecked-get js/process "env") "ARDZ_NET_DEBUG"))
(js/console.debug
"[net] inbound" (.slice id -6) "skipped:"
(cond
(identical? id (unchecked-get self "_myId")) "self"
(identical? id (unchecked-get self "_relayId")) "relay"
(.has (env/self-peer-ids) id) "sibling"
(.has ^js (unchecked-get self "_rejected") id) "rejected"
:else "wrong-room")))
(.abort ^js stream (js/Error. "unwelcome")))
(pm-adopt self (pm-state self (.-remotePeer ^js conn)) stream)))
js/undefined)
(defn- collision-settled?
"A second stream for the same peer: keep the one initiated by the LOWER peer
id. Returns false when `stream` loses and was aborted (the TS original's
early return), true when adoption should proceed."
[self p stream]
(let [cur (unchecked-get p "stream")]
(if-not (and ^boolean (js* "!!(~{})" cur) (not (identical? cur stream)))
true
(let [preferred (if (< (unchecked-get self "_myId") (unchecked-get p "id")) "outbound" "inbound")]
(if (and (identical? (.-direction ^js cur) preferred)
(identical? (.-status ^js cur) "open"))
(do (.abort ^js stream (js/Error. "duplicate stream")) false)
(do (.abort ^js cur (js/Error. "superseded")) true))))))
(defn- pm-adopt [self p stream]
(js/console.debug "[net] adopt stream" (.slice (unchecked-get p "id") -6)
(.-direction ^js stream) (.-status ^js stream))
(when (collision-settled? self p stream)
(when-not (identical? 0 (unchecked-get p "fallbackTimer"))
(js/clearTimeout (unchecked-get p "fallbackTimer"))
(unchecked-set p "fallbackTimer" 0))
(when-not (identical? 0 (unchecked-get p "redialTimer"))
(js/clearTimeout (unchecked-get p "redialTimer"))
(unchecked-set p "redialTimer" 0))
(unchecked-set p "stream" stream)
(unchecked-set p "verified" false)
(let [decoder (LengthPrefixedDecoder. (js-obj "maxDataLength" MAX-FRAME-BYTES))]
(unchecked-set p "frameChain" (js/Promise.resolve nil))
(.addEventListener
^js stream "message"
(fn [evt]
(when (identical? (unchecked-get p "stream") stream)
(try
(let [chunks (.decode ^js decoder (.-data ^js evt))]
(doseq [chunk (array-seq (js/Array.from chunks))]
(let [text (.decode td (.subarray ^js chunk))]
(unchecked-set
p "frameChain"
(.catch (.then (unchecked-get p "frameChain")
(fn [_] (pm-on-frame self p stream text)))
(fn [_] nil))))))
(catch :default _
(.abort ^js stream (js/Error. "bad framing")))))))
(.addEventListener ^js stream "close"
(fn [] (pm-on-stream-closed self p stream))))
;; both sides prove membership immediately (identity rides along; old
;; clients ignore the extra field)
(let [hello-id (unchecked-get self "_helloId")
mine (when (some? hello-id) (unchecked-get hello-id "mine"))
extra (if ^boolean (js* "!!(~{})" mine) (js-obj "id" mine) js/undefined)]
(-> (rc/seal-hello (unchecked-get self "_rc") (unchecked-get self "_myId")
(unchecked-get p "id") extra)
(.then (fn [env] (pm-write self stream env)))
(.catch (fn [_] (.abort ^js stream (js/Error. "seal failed"))))))
(unchecked-set
p "helloTimer"
(js/setTimeout
(fn []
(when (and (identical? (unchecked-get p "stream") stream)
(not (unchecked-get p "verified")))
(.abort ^js stream (js/Error. "hello timeout"))))
HELLO-TIMEOUT-MS)))
js/undefined)
(defn- pm-on-frame [self p stream text]
(-> (js/Promise.resolve nil)
(.then
(fn [_]
(let [parsed (try (js/JSON.parse text) (catch :default _ ::bad))]
(cond
(identical? parsed ::bad) nil
(unchecked-get p "verified")
(do (ev-call self "message" (unchecked-get p "id") parsed) nil)
:else
;; first frame must be the remote's sealed hello
(js-await [pt (rc/open-hello (unchecked-get self "_rc") (unchecked-get p "id")
(unchecked-get self "_myId") parsed)]
(when (identical? (unchecked-get p "stream") stream) ; not superseded
(if (nil? pt)
(do
(when ^boolean (js* "!!(~{})"
(unchecked-get (unchecked-get js/process "env") "ARDZ_NET_DEBUG"))
(js/console.debug "[net] bad hello from" (.slice (unchecked-get p "id") -6)
(.-direction ^js stream) "frame:" (.slice text 0 100)))
(.add ^js (unchecked-get self "_rejected") (unchecked-get p "id"))
(.abort ^js stream (js/Error. "bad hello"))
(if (unchecked-get p "member")
(pm-drop self p)
(.delete ^js (unchecked-get self "_peers") (unchecked-get p "id")))
nil)
(do
(unchecked-set p "verified" true)
(unchecked-set p "redialAttempt" 0)
(when-not (identical? 0 (unchecked-get p "helloTimer"))
(js/clearTimeout (unchecked-get p "helloTimer"))
(unchecked-set p "helloTimer" 0))
(when-not (unchecked-get p "member")
(unchecked-set p "member" true)
(ev-call self "peerOpen" (unchecked-get p "id"))
(pm-emit-members self))
;; identity is additive: membership never depends on it,
;; and verification happens off the hot path (fires
;; peerIdentity when it lands)
(let [hello-id (unchecked-get self "_helloId")]
(when (and ^boolean (js* "!!(~{})" hello-id)
(not (identical? js/undefined (unchecked-get pt "id")))
(not ^boolean (js* "!!(~{})" (unchecked-get p "ident"))))
(-> ((unchecked-get hello-id "verify") (unchecked-get p "id")
(unchecked-get pt "id"))
(.then (fn [ident]
(when (and ^boolean (js* "!!(~{})" ident)
(identical?
(.get ^js (unchecked-get self "_peers")
(unchecked-get p "id"))
p))
(unchecked-set p "ident" ident)
(ev-call self "peerIdentity" (unchecked-get p "id") ident)))))))
nil))))))))))
(defn- pm-on-stream-closed [self p stream]
(when (identical? (unchecked-get p "stream") stream)
(unchecked-set p "stream" nil)
(unchecked-set p "verified" false)
(when-not (identical? 0 (unchecked-get p "helloTimer"))
(js/clearTimeout (unchecked-get p "helloTimer"))
(unchecked-set p "helloTimer" 0))
(if (unchecked-get p "member")
(pm-schedule-redial self p)
(.delete ^js (unchecked-get self "_peers") (unchecked-get p "id"))))
js/undefined)
(defn- pm-schedule-redial [self p]
(when (identical? 0 (unchecked-get p "redialTimer"))
(let [delay (aget REDIAL-DELAYS-MS (unchecked-get p "redialAttempt"))]
(if (identical? delay js/undefined)
(pm-drop self p)
(do
(unchecked-set p "redialAttempt" (inc (unchecked-get p "redialAttempt")))
(unchecked-set
p "redialTimer"
(js/setTimeout
(fn []
(unchecked-set p "redialTimer" 0)
(when (and (unchecked-get p "member")
(not ^boolean (js* "!!(~{})" (unchecked-get p "stream")))
(not (unchecked-get p "dialing")))
(pm-dial self p)))
delay))))))
js/undefined)
(defn- pm-drop [self p]
(doseq [t [(unchecked-get p "helloTimer")
(unchecked-get p "fallbackTimer")
(unchecked-get p "redialTimer")]]
(when-not (identical? 0 t) (js/clearTimeout t)))
(let [s (unchecked-get p "stream")]
(when ^boolean (js* "!!(~{})" s) (.abort ^js s (js/Error. "dropped"))))
(.delete ^js (unchecked-get self "_peers") (unchecked-get p "id"))
(when (unchecked-get p "member")
(unchecked-set p "member" false)
(ev-call self "peerGone" (unchecked-get p "id"))
(pm-emit-members self))
js/undefined)
(defn- pm-emit-members [self]
(ev-call self "membersChanged" (members self)))
(defn start [self]
(-> (js/Promise.resolve nil)
(.then
(fn [_]
(js-await [_ (.handle ^js (unchecked-get self "_node")
(unchecked-get self "protocol")
(fn [stream conn] (pm-on-inbound self stream conn))
(js-obj "runOnLimitedConnection" true))]
(do
(.addEventListener ^js (unchecked-get self "_node") "peer:discovery"
(fn [e] (pm-on-discovered self (unchecked-get (.-detail ^js e) "id"))))
js/undefined))))))
(let [proto (.-prototype PeerManager)]
(unchecked-set proto "start" (fn [] (this-as self (start self))))
(unchecked-set proto "members" (fn [] (this-as self (members self))))
(unchecked-set proto "openPeers" (fn [] (this-as self (open-peers self))))
(unchecked-set proto "openCount" (fn [] (this-as self (open-count self))))
(unchecked-set proto "identityOf" (fn [id] (this-as self (identity-of self id))))
(unchecked-set proto "sendTo" (fn [id frame] (this-as self (send-to self id frame))))
(unchecked-set proto "broadcast" (fn [frame] (this-as self (broadcast self frame))))
(unchecked-set proto "wake" (fn [] (this-as self (wake self)))))
|