chat / client / src / sueta / lib / mailbox.cljs
   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
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
;; ported-from: src/lib/mailbox.ts — and, since Phase 6b, no longer a
;; transliteration of it.
;;
;; Offline delivery via the relay node's store-and-forward mailbox
;; (PROTOCOL.md §9). The mailbox is a second transport for the SAME bytes the
;; log replicates — sealed entry blocks, unixfs image blocks and (sealed for
;; transit) identity record blocks — so two members who are never online
;; together still converge. The log stays the single source of truth; the node
;; stores opaque blobs it cannot read, deduped per room by SHA-256.
;;
;; Deposit is author-only (each device uploads only what it created; the
;; persistent queue IS the catch-up story, server dedup absorbs re-sends).
;; Replay is two-phase: store every block first, then join in seq order, so no
;; join ever waits on a bitswap timeout. Unjoinable entries (ancestors expired
;; from the ring buffer) surface display-only and sit on a retry list drained
;; when live replication backfills the gap.
;;
;; Everything here is best-effort: a mailbox failure must never break live chat,
;; so every path ends in a typed result or a logged swallow.
;;
;; ── WHY THIS FILE IS WRITTEN THE WAY IT IS ───────────────────────────────────
;;
;; This is the file that took chat down twice, and both defects were the same
;; defect: a `js-await` ladder deep enough that a paren landing one level early
;; still LOOKED right. `do-replay` was seventy lines and nine levels of nesting;
;; the trailing `.catch` and `.then` of its outer chain ended up as extra
;; ARGUMENTS to the `.then` above them, compiled into a method call on a
;; function literal, and threw the instant boot called it. Rooms would not open
;; at all, and — twelve lines away, same cause — offline delivery had never once
;; worked. Both suites were green either side of both fixes.
;;
;; So the shape below is the durable fix, not a tidy-up:
;;
;;  · `p/let` instead of `js-await`. `shadow.cljs.modern/js-await` is pure
;;    `.then` sugar, so every await nests one level deeper than the last; `p/let`
;;    is the identical chain with the nesting spent in a binding vector. Nothing
;;    here is more than three levels deep any more, and the longest closing run
;;    is five parens.
;;  · `p/let` SEQUENCES, and that is exactly why it is safe here. PROTOCOL.md §9
;;    fixes the side-effect ORDER of replay and deposit — verify and store every
;;    page, THEN join in seq order — so `p/all` is not an optimisation available
;;    to a later reader. Neither is `j/each-in-order!` becoming `Promise.all`.
;;  · every promise-returning step is one named function with one job, so the
;;    order of the effects is legible as a list of names rather than as
;;    indentation. `test/vectors/mailbox-effects.json` asserts that order.
;;  · a `p/let` BODY awaits every expression in it, one microtask each — so the
;;    synchronous tails below are each ONE call to a named `!`-suffixed
;;    function (`deposit-step!`, `collect-entry!`, `finish-replay!`) whose body
;;    is an ordinary `do`. Spreading those statements across a multi-expression
;;    body would interleave microtasks between lines whose order is the
;;    protocol. Same rule as sueta.room's.
;;
;; ── WHAT IS CONSTRAINED HERE, AND BY WHAT ────────────────────────────────────
;;
;;  · the DEPOSIT ENVELOPES are wire bytes and the QUEUE ITEMS are JSON in
;;    localStorage, so both are key-ordered contracts. Queue items are declared
;;    once, in QUEUE-ITEM below, and built through sueta.wire/encode — the same
;;    sequential unchecked-set `j/ordered` does, so the bytes do not move;
;;    test/vectors/mailbox.json pins all three shapes and both envelopes.
;;  · the CID RE-HASH before storage is security-critical (`store-block!`).
;;  · the CURSOR is an opaque server string compared by (length, lexicographic)
;;    and must never be numeric-parsed — `seq-lt` is that rule.
;;  · JS truthiness where the TypeScript leaned on it: a `_retryTimer` of 0
;;    means "no timer", `_lastReplay` of 0 means "never replayed", the persisted
;;    ident flag is the string "1", and `bytes` may be a zero-length array.
;;  · the twelve `_`-prefixed methods at the bottom are a test seam
;;    (test/vectors/facade.json freezes the set, the order and the arities).
(ns sueta.lib.mailbox
  (:require ["multiformats/cid" :refer (CID)]
            ["multiformats/bases/base58" :refer (base58btc)]
            [promesa.core :as p]
            [sueta.fx :as fx]
            [sueta.wire :as w]
            [ardegazu.rooms.js :as j]
            [ardegazu.rooms.lib.crypto :as rc]
            [shadow.cljs.modern :refer (defclass)]))

(def ^:private MINT-TTL-SEC 86400) ; ask for the max (24 h); the node may clamp
(def ^:private REPLAY-THROTTLE-MS 30000)
(def ^:private BACKOFF-BASE-MS 1000)
(def ^:private BACKOFF-CAP-MS 300000)
;; There is deliberately NO cap on the outgoing queue. It holds THIS user's own
;; not-yet-deposited ops — bounded by real usage, not by adversarial input — and
;; the suite's standing rule is: bound WORK, never HISTORY; pruning is a user
;; action, never a silent one. The cap this replaces was a drop-oldest at 500,
;; which silently discarded the EARLIEST queued item (for a payments app, a
;; payment order). What stays bounded is the work: one deposit in flight
;; (`flush!`), one page walk per replay, and the SENT ring below. If persisting
;; the queue fails, the failure is SURFACED through onNotice — see save-queue!.
;;
;; SENT-CAP verdict: this one is a dedup/bookkeeping RING, and a cap is correct
;; here. `_sent` only saves round trips — it short-circuits re-enqueueing keys
;; already deposited. Evicting an old key loses no data and re-sends nothing by
;; itself; at worst a later re-enqueue of an evicted key costs one redundant
;; deposit, which the node's per-room SHA-256 dedup collapses (PROTOCOL.md §9).
;; Bounding it bounds the `.includes` scan and the localStorage record: work,
;; not history.
(def ^:private SENT-CAP 1000)
(def ^:private MAX-REPLAY-PAGES 200)
(def ^:private DAG-CBOR 0x71)
(def ^:private SHA2-256 0x12)

;; Deposit envelope, one mailbox message per block (deterministic so the node's
;; SHA-256 dedup collapses re-sends):
;;   0x01 | cid.bytes | block bytes          raw blocks (sealed entries, image blocks)
;;   0x02 | seal(cid.bytes | block bytes)    identity records (plaintext blocks — §7b)
(def ^:private RAW-BLOCK 0x01)
(def ^:private SEALED-BLOCK 0x02)
(def ^:private RAW-PREFIX (js/Uint8Array. #js [RAW-BLOCK]))
(def ^:private SEALED-PREFIX (js/Uint8Array. #js [SEALED-BLOCK]))

;; ---- pure helpers ----------------------------------------------------------
;;
;; Nothing below this heading touches `self`, the clock, the network or storage.

(defn- concat-bytes
  "A fresh Uint8Array of every part, in order. The envelope builder."
  [parts]
  (let [out (js/Uint8Array. (reduce (fn [n p] (+ n (.-length ^js p))) 0 parts))]
    (reduce (fn [off p] (.set out p off) (+ off (.-length ^js p))) 0 parts)
    out))

(defn- ^boolean bytes-eq
  "Length-then-XOR compare, no early exit — this decides whether a replayed
   block matches the CID it arrived under, so it must not leak where two
   digests first differ."
  [a b]
  (and (identical? (.-length ^js a) (.-length ^js b))
       (loop [i 0 diff 0]
         (if (>= i (.-length ^js a))
           (zero? diff)
           (recur (inc i) (bit-or diff (bit-xor (aget a i) (aget b i))))))))

(defn- ^boolean seq-lt
  "a < b for mailbox seqs — never numeric-parse them (they overflow doubles);
   compare by (length, lexicographic).

   `(< a b)` on STRINGS is correct here: CLJS's two-argument `<` inlines to the
   raw JS `<` operator, so it is JS string comparison, not a numeric coercion.
   (The same inlining is what makes the reaction fold's `(> hash-a hash-b)`
   tiebreak work in lib/log and app/chat.)"
  [a b]
  (if-not (identical? (.-length ^string a) (.-length ^string b))
    (< (.-length ^string a) (.-length ^string b))
    (< a b)))

(defn- ^boolean unauthorized?
  "The two statuses that mean 'your bearer is stale', which is the only status
   pair worth a re-mint."
  [status]
  (or (identical? 401 status) (identical? 403 status)))

(defn- backoff-ms
  "The nth retry delay: 1s, 2s, 4s … capped at 5 minutes. The ladder
   test/vectors/mailbox-effects.json pins IS this function over 0,1,2,…"
  [n]
  (min (* BACKOFF-BASE-MS (js/Math.pow 2 n)) BACKOFF-CAP-MS))

;; ---- the storage boundary --------------------------------------------------
;;
;; A queue item is JSON in localStorage, so its key ORDER is a contract with
;; every browser that already holds one. Declared here once instead of at three
;; call sites; `w/encode` sets the fields in declaration order.

(def ^:private QUEUE-ITEM
  {"entry" [[:t "t"] [:hash "hash"]]
   "img" [[:t "t"] [:cid "cid"] [:group "group"]]
   "ident" [[:t "t"]]})

(defn- item
  "One persisted queue item of kind `t`."
  [t m]
  (w/encode (QUEUE-ITEM t) (assoc m :t t)))

(defn- item-t [it] (unchecked-get it "t"))
(defn- item-hash [it] (unchecked-get it "hash"))
(defn- item-cid [it] (unchecked-get it "cid"))
(defn- item-group [it] (unchecked-get it "group"))

;; localStorage, tolerant of Safari private mode / quota (mailbox state is
;; best-effort: losing it costs re-deposits the server dedups away)
(defn- load-str [k]
  (try (js/localStorage.getItem k) (catch :default _ nil)))

(defn- save-str [k v]
  (try (js/localStorage.setItem k v) (catch :default _ nil))
  js/undefined)

(defn- load-json [k fallback]
  (try
    (let [raw (js/localStorage.getItem k)]
      (if (j/truthy? raw) (js/JSON.parse raw) fallback))
    (catch :default _ fallback)))

(defn- ^boolean save-json
  "True when the write landed. Most callers ignore the verdict (sent/retry are
   best-effort bookkeeping); save-queue! does not — see it for why."
  [k v]
  (try (js/localStorage.setItem k (js/JSON.stringify v)) true
       (catch :default _ false)))

;; ---- auth: browser open-mint, cached bearer (turn.cljs pattern) ------------

(defclass MailboxAuth
  (constructor [this creds-url room-id]
    (unchecked-set this "_credsUrl" creds-url)
    (unchecked-set this "_roomId" room-id)
    (unchecked-set this "_cached" nil)
    (unchecked-set this "_validUntil" 0)
    ;; an atom, not a field, because fx/single-flight! is what names this
    ;; pattern and it needs a cell it can clear from a `.finally`
    (unchecked-set this "_inflight" (atom nil))))

(defn- cache-token!
  "Hold the bearer, and renew five minutes before the node says it dies."
  [self token ttl-sec]
  (unchecked-set self "_cached" token)
  (unchecked-set self "_validUntil" (- (+ (js/Date.now) (* ttl-sec 1000)) 300000))
  token)

(defn- mint
  "One open mint against the relay. The node clamps the ttl it feels like
   giving; we renew five minutes before whatever it said."
  [self]
  (let [url (str (unchecked-get self "_credsUrl")
                 "?room=" (js/encodeURIComponent (unchecked-get self "_roomId"))
                 "&ttl=" MINT-TTL-SEC)]
    (p/let [res (js/fetch url (j/ordered "signal" (js/AbortSignal.timeout 15000)))
            _ (when-not (j/truthy? (unchecked-get res "ok"))
                (throw (js/Error. (str "mailbox mint " (unchecked-get res "status")))))
            body (.json ^js res)]
      (let [token (unchecked-get body "token")]
        (when-not (j/truthy? token)
          (throw (js/Error. "mailbox mint: no token")))
        (cache-token! self token (j/nn (unchecked-get body "ttl") 0))))))

(defn auth-get
  "The cached bearer while it is fresh, otherwise one mint shared by everyone
   who asks while it is running."
  [self]
  (if (and (j/truthy? (unchecked-get self "_cached"))
           (< (js/Date.now) (unchecked-get self "_validUntil")))
    (p/resolved (unchecked-get self "_cached"))
    (fx/single-flight! (unchecked-get self "_inflight") #(mint self))))

(defn auth-invalidate
  "Drop the cached token (server said 401/403) so the next get() re-mints."
  [self]
  (unchecked-set self "_cached" nil)
  (unchecked-set self "_validUntil" 0)
  js/undefined)

(let [proto (.-prototype MailboxAuth)]
  (unchecked-set proto "get" (fn [] (this-as self (auth-get self))))
  (unchecked-set proto "invalidate" (fn [] (this-as self (auth-invalidate self))))
  (unchecked-set proto "_mint" (fn [] (this-as self (mint self)))))

;; ---- bare HTTP client with typed results ------------------------------------

(defclass MailboxClient
  (constructor [this base-url room-id auth]
    (unchecked-set this "_auth" auth)
    (unchecked-set this "_base"
                   (str (.replace ^string base-url (js/RegExp. "\\/$") "")
                        "/rooms/" (js/encodeURIComponent room-id) "/messages"))))

(defn- authed-init
  "`(init-fn)` plus the bearer and a 20 s abort — a fresh options bag, so the
   caller's constructor never sees a mutation."
  [base token]
  (let [opts (js/Object.assign (js-obj) base)
        headers (js/Object.assign (js-obj) (j/nn (unchecked-get base "headers") (js-obj)))]
    (unchecked-set headers "authorization" (str "Bearer " token))
    (unchecked-set opts "headers" headers)
    (unchecked-set opts "signal" (js/AbortSignal.timeout 20000))
    opts))

(defn- send-req
  "One request with the cached bearer; on 401/403 re-mint ONCE and retry.
   Resolves to the Response, or to nil for every failure a caller is expected
   to classify as `net` — no token, a refused fetch, a second 401.

   The two `nil`s are deliberately indistinguishable: PROTOCOL.md §9 makes the
   mailbox best-effort, and a caller that could tell 'no token' from 'no
   network' would only be tempted to act differently on them."
  [self init-fn qs]
  (letfn [(attempt [n]
            (if (>= n 2)
              (p/resolved nil)
              (p/let [token (fx/drop-silently (j/later #(auth-get (unchecked-get self "_auth"))))]
                (when (some? token)
                  (p/let [res (fx/drop-silently
                               (j/later #(js/fetch (str (unchecked-get self "_base") qs)
                                                   (authed-init (init-fn) token))))]
                    (if (and (some? res)
                             (zero? n)
                             (unauthorized? (unchecked-get res "status")))
                      (do (auth-invalidate (unchecked-get self "_auth")) (attempt (inc n)))
                      res))))))]
    (attempt 0)))

(defn deposit
  "POST one envelope. Every outcome is a typed record, never a throw: `ok` with
   the node's dedup verdict and seq, or `kind` in
   too_large / auth / backpressure / net."
  [self payload]
  (p/let [res (send-req self
                        (fn [] (j/ordered "method" "POST"
                                          "headers" (j/ordered "content-type" "application/octet-stream")
                                          "body" payload))
                        "")]
    (if (nil? res)
      (j/ordered "ok" false "kind" "net")
      (let [status (unchecked-get res "status")]
        (cond
          (or (identical? 201 status) (identical? 200 status))
          ;; a 2xx with an unreadable body still deposited the bytes
          (-> (.json ^js res)
              (.then (fn [b] (j/ordered "ok" true
                                        "deduped" (j/truthy? (unchecked-get b "deduped"))
                                        "seq" (j/nn (unchecked-get b "seq") ""))))
              (.catch (fn [_] (j/ordered "ok" true "deduped" false "seq" ""))))

          (identical? 413 status) (j/ordered "ok" false "kind" "too_large")
          (unauthorized? status) (j/ordered "ok" false "kind" "auth")
          :else (j/ordered "ok" false "kind"
                           (if (or (identical? 429 status) (identical? 507 status)) "backpressure" "net")))))))

(defn replay-page
  "GET one page after `after` (nil = from the beginning). `messages` is clamped
   to an array — `Array.isArray`, not `array?`: a validator IS the wire
   contract, and `array?` is realm-sensitive."
  [self after limit]
  (let [qs (str "?limit=" limit (if (j/truthy? after) (str "&after=" (js/encodeURIComponent after)) ""))]
    (p/let [res (send-req self (fn [] (j/ordered "method" "GET")) qs)]
      (if (nil? res)
        (j/ordered "ok" false "kind" "net")
        (let [status (unchecked-get res "status")]
          (cond
            (unauthorized? status) (j/ordered "ok" false "kind" "auth")

            (not (j/truthy? (unchecked-get res "ok")))
            (j/ordered "ok" false "kind" (if (identical? 429 status) "backpressure" "net"))

            :else
            (-> (.json ^js res)
                (.then (fn [b]
                         (let [msgs (unchecked-get b "messages")]
                           (j/ordered
                            "ok" true
                            "page" (j/ordered "messages" (if (w/arr? msgs) msgs (array))
                                              "hasMore" (j/truthy? (unchecked-get b "has_more"))
                                              "oldestSeq" (j/nn (unchecked-get b "oldest_seq") nil)
                                              "latestSeq" (j/nn (unchecked-get b "latest_seq") nil))))))
                (.catch (fn [_] (j/ordered "ok" false "kind" "net"))))))))))

(let [proto (.-prototype MailboxClient)]
  (unchecked-set proto "deposit" (fn [payload] (this-as self (deposit self payload))))
  (unchecked-set proto "replay" (fn [after limit] (this-as self (replay-page self after limit))))
  (unchecked-set proto "_send" (fn [init-fn qs] (this-as self (send-req self init-fn qs)))))

;; ---- the orchestrator --------------------------------------------------------

(defclass MailboxSync
  (constructor [this opts]
    (unchecked-set this "_log" (unchecked-get opts "log"))
    (unchecked-set this "_client" (unchecked-get opts "client"))
    (unchecked-set this "_cipher" (unchecked-get opts "cipher"))
    (unchecked-set this "_maxBytes" (* (unchecked-get opts "maxMessageKb") 1024))
    (unchecked-set this "_keys" (unchecked-get opts "keys"))
    (unchecked-set this "_hasIdentity" (unchecked-get opts "hasIdentity"))
    (unchecked-set this "_onDisplayOnly" (j/nn (unchecked-get opts "onDisplayOnly") (fn [_] js/undefined)))
    (unchecked-set this "_onNotice" (j/nn (unchecked-get opts "onNotice") (fn [_] js/undefined)))
    (unchecked-set this "_onReplayed" (j/nn (unchecked-get opts "onReplayed") (fn [_] js/undefined)))
    (unchecked-set this "_flushing" false)
    (unchecked-set this "_replaying" false)
    (unchecked-set this "_lastReplay" 0)
    (unchecked-set this "_attempt" 0)
    (unchecked-set this "_retryTimer" nil)
    (unchecked-set this "_liveTimer" nil)
    (unchecked-set this "_noticed" false)
    (unchecked-set this "_qNoticed" false)
    ;; the mutex that keeps image blocks ahead of their entry (fx/locked!)
    (unchecked-set this "_enqueueChain" (js/Promise.resolve nil))
    (unchecked-set this "_listeners" (array))
    (let [ks (unchecked-get opts "keys")]
      (unchecked-set this "_queue" (load-json (unchecked-get ks "queue") (array)))
      ;; hashes/CIDs already deposited (ring; server dedup is the real guard)
      (unchecked-set this "_sent" (load-json (unchecked-get ks "sent") (array)))
      ;; stored-but-unjoinable entry hashes
      (unchecked-set this "_retry" (load-json (unchecked-get ks "retry") (array))))))

(declare flush! do-replay drain-retry join-phase process-message build-payload
         enqueue-identity-once push-item drop-head mark-sent schedule-retry
         add-retry remove-retry)

;; ── the JS boundary for this class ───────────────────────────────────────────
;;
;; MailboxSync's state is a string-keyed JS object (`_attempt` is read by
;; test/helpers/mailbox-effects.mjs, so this is a seam, not a preference).
;; Everything that reads it does so through the names below, so a field name is
;; a string literal in exactly one place and the steps above read as prose.
;;
;; `^js` goes on the accessor's VAR, and again on the call form at any interop
;; site, because a type hint is lost through macroexpansion — `^js
;; (unchecked-get self "_log")` reads as tagged and is not (dev/docs/CLJS.md).
;; externs.js exists because of exactly that.

(defn- key-of [self k] (unchecked-get (unchecked-get self "_keys") k))
(defn- ^js log-of [self] (unchecked-get self "_log"))
(defn- ^js client-of [self] (unchecked-get self "_client"))
(defn- queue-of [self] (unchecked-get self "_queue"))
(defn- sent-of [self] (unchecked-get self "_sent"))
(defn- retry-of [self] (unchecked-get self "_retry"))

;; Timer ids are read with JS truthiness on purpose: an id of 0 is a real id in
;; some hosts and `0` is truthy in Clojure, so `armed?` is the only safe test.
(defn- ^boolean retry-armed? [self] (j/truthy? (unchecked-get self "_retryTimer")))
(defn- ^boolean live-armed? [self] (j/truthy? (unchecked-get self "_liveTimer")))

(defn- clear-retry! [self]
  (when (retry-armed? self) (js/clearTimeout (unchecked-get self "_retryTimer")))
  (unchecked-set self "_retryTimer" nil)
  js/undefined)

(defn- clear-live! [self]
  (when (live-armed? self) (js/clearTimeout (unchecked-get self "_liveTimer")))
  (unchecked-set self "_liveTimer" nil)
  js/undefined)

;; The three callbacks the room installs. Defaulted in the constructor, so no
;; call site ever has to test for one.
(defn- on-replayed! [self n] ((unchecked-get self "_onReplayed") n) js/undefined)
(defn- on-notice! [self msg] ((unchecked-get self "_onNotice") msg) js/undefined)
(defn- on-display-only! [self meta-obj] ((unchecked-get self "_onDisplayOnly") meta-obj) js/undefined)

(defn- save-queue!
  "Persist the queue — and, because the queue is uncapped and holds the user's
   own undeposited ops, a persist that FAILS (quota, private mode) is surfaced
   through onNotice, once per session, instead of being swallowed: the items
   still deposit from memory this session, but will not survive a page close."
  [self]
  (when-not (save-json (key-of self "queue") (queue-of self))
    (when-not (j/truthy? (unchecked-get self "_qNoticed"))
      (unchecked-set self "_qNoticed" true)
      (on-notice! self "this browser is not persisting queued messages — they may be lost if you close this page")))
  js/undefined)

;; ── decoding MailboxClient's typed records ───────────────────────────────────
;;
;; `deposit` and `replay` answer with the JS records above — a shape a test
;; double also produces, so this is a real boundary and not an internal call.
;; The two readers below are the only place their fields are named, and the only
;; place JS truthiness is applied to them; above this line the results are
;; Clojure maps and `nil` is the only absence.

(defn- deposit-result [r]
  {:ok? (j/truthy? (w/oget r "ok"))
   :kind (w/oget r "kind")})

(defn- replay-result [r]
  (let [pg (w/oget r "page")]
    {:ok? (j/truthy? (w/oget r "ok"))
     ;; `messages` stays a JS array: replay-page already clamped it with
     ;; Array.isArray, and it is walked by index
     :messages (w/oget pg "messages")
     :more? (j/truthy? (w/oget pg "hasMore"))
     :oldest-seq (w/non-empty (w/oget pg "oldestSeq"))}))

(defn start
  "Non-blocking: replay (which also drains the retry list), then flush. The
   wake kick resets the backoff and cancels a pending retry — coming back from
   a locked screen is new information about the network, not another failure."
  [self]
  (let [kick (fn []
               (when-not (identical? "hidden" (unchecked-get js/document "visibilityState"))
                 (unchecked-set self "_attempt" 0)
                 (clear-retry! self)
                 (flush! self)
                 (do-replay self false))
               js/undefined)
        listeners #js [#js ["visibilitychange" kick] #js ["pageshow" kick] #js ["online" kick]]]
    (unchecked-set self "_listeners" listeners)
    (doseq [pair (array-seq listeners)]
      (.addEventListener js/window (aget pair 0) (aget pair 1)))
    (-> (do-replay self true)
        (.then (fn [_] (flush! self))))
    js/undefined))

(defn stop [self]
  (doseq [pair (array-seq (unchecked-get self "_listeners"))]
    (.removeEventListener js/window (aget pair 0) (aget pair 1)))
  (clear-retry! self)
  (clear-live! self)
  js/undefined)

(defn state
  "Debug/e2e visibility."
  [self]
  (j/ordered "queued" (.-length (queue-of self))
             "retry" (.-length (retry-of self))
             "cursor" (load-str (key-of self "cursor"))))

;; ---- deposit side ----------------------------------------------------------

(defn enqueue-entry
  "RoomLog.onAppended — this device authored a new entry. Synchronous up to the
   persisted queue, so a page closed right after send still deposits on the next
   visit."
  [self hash]
  (when-not (.includes (sent-of self) hash)
    (enqueue-identity-once self)
    (push-item self (item "entry" {:hash hash}))
    (save-queue! self)
    (flush! self))
  js/undefined)

(defn- enqueue-blocks!
  "Children first, root last — the order `imageDagBlocks` returns, and the order
   a receiver needs to walk the DAG without a fetch."
  [self blocks root-cid]
  (enqueue-identity-once self)
  (doseq [b (array-seq blocks)]
    (let [cid (.toString (unchecked-get b "cid"))]
      (when-not (.includes (sent-of self) cid)
        (push-item self (item "img" {:cid cid :group root-cid})))))
  (save-queue! self)
  js/undefined)

(defn enqueue-image
  "RoomLog.onImageStored — expand the image DAG into per-block items.

   Serialised through `_enqueueChain` so two images in flight cannot interleave
   their expansions into the queue. The expansion is async, so image blocks may
   still enqueue AFTER their entry; that is fine, because replay stores every
   block before it joins anything, so order across items never matters."
  [self root-cid]
  (fx/locked!
   self "_enqueueChain"
   (fn []
     (-> (fx/tolerate
          "mailbox image enqueue failed"
          (p/let [blocks (.imageDagBlocks ^js (log-of self) root-cid)]
            (enqueue-blocks! self blocks root-cid)))
         (.then (fn [_] (flush! self) js/undefined)))))
  js/undefined)

(defn- enqueue-identity-once
  "Queue our identity record, once per room and AHEAD of any entry — receivers
   need it for canAppend. The caller persists."
  [self]
  (when (and (j/truthy? (unchecked-get self "_hasIdentity"))
             (not (j/truthy? (load-str (key-of self "ident"))))
             (not (.some (queue-of self) (fn [i] (identical? "ident" (item-t i))))))
    (.push (queue-of self) (item "ident" {})))
  js/undefined)

(defn- push-item
  "Append — and nothing else. This used to drop-oldest past a cap of 500, which
   silently discarded the user's own earliest undeposited op; the queue is now
   uncapped (see the policy note by SENT-CAP). Kept as a function because
   `_push` is a frozen test seam."
  [self it]
  (.push (queue-of self) it)
  js/undefined)

(defn- drop-head
  "Drop the head item — and, for an image block, its whole DAG group: half an
   image is not worth the round trips."
  [self it]
  (.shift (queue-of self))
  (when (identical? "img" (item-t it))
    (let [group (item-group it)]
      (unchecked-set self "_queue"
                     (.filter (queue-of self)
                              (fn [i] (not (and (identical? "img" (item-t i))
                                                (identical? group (item-group i)))))))))
  (save-queue! self)
  js/undefined)

(defn- mark-sent [self k]
  (let [sent (sent-of self)]
    (.push sent k)
    (when (> (.-length sent) SENT-CAP)
      (.splice sent 0 (- (.-length sent) SENT-CAP)))
    (save-json (key-of self "sent") sent))
  js/undefined)

(defn- schedule-retry
  "Arm the backoff timer, unless one is already armed. `_retryTimer` is read
   with JS truthiness because a timer id of 0 means 'no timer'."
  [self]
  (when-not (retry-armed? self)
    (let [n (unchecked-get self "_attempt")]
      (unchecked-set self "_attempt" (inc n))
      (unchecked-set self "_retryTimer"
                     (js/setTimeout (fn []
                                      (unchecked-set self "_retryTimer" nil)
                                      (flush! self)
                                      js/undefined)
                                    (backoff-ms n)))))
  js/undefined)

(defn- build-payload
  "The deposit envelope for one queue item, or nil when it cannot be built — a
   block that is gone, a malformed cid, no identity record yet. nil means DROP,
   handled at the call site exactly as a 413 is."
  [self it]
  (let [^js log (log-of self)]
    (case (item-t it)
      "entry"
      (p/let [bytes (.sealedEntryBytes log (item-hash it))]
        (when (j/truthy? bytes)
          ;; CID.parse throws on a hash that is not base58btc dag-cbor
          (try (concat-bytes [RAW-PREFIX
                              (unchecked-get (.parse CID (item-hash it) base58btc) "bytes")
                              bytes])
               (catch :default _ nil))))

      "img"
      (fx/drop-silently
       (p/let [cid (.parse CID (item-cid it))
               bytes (.rawBlock log cid)]
         (when (j/truthy? bytes)
           (concat-bytes [RAW-PREFIX (unchecked-get cid "bytes") bytes]))))

      (let [blk (.myIdentityBlock log)]
        (if-not (j/truthy? blk)
          (p/resolved nil)
          (p/let [sealed (rc/at-rest-seal (unchecked-get self "_cipher")
                                          (concat-bytes [(unchecked-get (unchecked-get blk "cid") "bytes")
                                                         (unchecked-get blk "bytes")]))]
            (concat-bytes [SEALED-PREFIX sealed])))))))

(defn- deposit-step!
  "One accepted deposit's bookkeeping, in ONE tick: the item leaves the queue,
   its key joins the sent ring (or the identity flag is written), the queue is
   persisted and the backoff ladder resets. Never split across awaits — a
   half-updated queue survives a page close."
  [self it]
  (.shift (queue-of self))
  (case (item-t it)
    "entry" (mark-sent self (item-hash it))
    "img" (mark-sent self (item-cid it))
    (save-str (key-of self "ident") "1"))
  (save-queue! self)
  (unchecked-set self "_attempt" 0)
  js/undefined)

(defn- end-flush!
  "Leave the critical section, then re-check: an enqueue that arrived while we
   were exiting saw `_flushing` and bailed, and its item must not sit stranded
   until the next wake kick."
  [self]
  (unchecked-set self "_flushing" false)
  (when (and (pos? (.-length (queue-of self))) (not (retry-armed? self)))
    (flush! self))
  js/undefined)

(defn- flush-one!
  "Build and deposit the head item, and answer whether the walk continues.

   Every branch that answers `true` has already removed the head, so the walk
   always makes progress; the one that answers `false` has armed the backoff
   timer, which is what will resume it."
  [self it]
  (p/let [payload (build-payload self it)]
    (if (or (nil? payload)
            (> (.-length ^js payload) (unchecked-get self "_maxBytes")))
      ;; unbuildable, or over the node's frame cap — dropped exactly like a 413
      (do (drop-head self it) true)
      (p/let [r (.deposit ^js (client-of self) payload)]
        (let [{:keys [ok? kind]} (deposit-result r)]
          (cond
            ok? (do (deposit-step! self it) true)
            (= "too_large" kind) (do (drop-head self it) true)
            ;; auth (the client already re-minted once), backpressure, net
            :else (do (schedule-retry self) false)))))))

(defn- flush!
  "Drain the queue oldest-first, ONE deposit in flight at a time.

   `_flushing` is the whole re-entrancy story: enqueue, the wake kick, the retry
   timer and `end-flush!` all call this, and exactly one of them walks.
   test/vectors/mailbox-effects.json's `flushReentrancy` asserts the
   consequence — two enqueues, two deposits, never a third and never concurrent
   — rather than the mechanism, which is the only way to test it honestly."
  [self]
  (if (or (j/truthy? (unchecked-get self "_flushing"))
          (zero? (.-length (queue-of self))))
    (p/resolved nil)
    (do
      (unchecked-set self "_flushing" true)
      (-> (letfn [(step []
                    (if (zero? (.-length (queue-of self)))
                      (p/resolved nil)
                      (p/let [go? (flush-one! self (aget (queue-of self) 0))]
                        (when go? (step)))))]
            (step))
          (.catch (fn [err]
                    ;; not silent: a flush that THROWS is our defect, and the
                    ;; retry ladder is what keeps the queue moving regardless
                    (js/console.error "mailbox flush failed" err)
                    (schedule-retry self)
                    nil))
          (.then (fn [_] (end-flush! self)))))))

;; ---- replay side -------------------------------------------------------------

(defn on-live-entry
  "Live replication delivered an entry — a gap may have closed; drain soon.
   Debounced to one drain per 1.5 s of arrivals, and `_liveTimer` is read with
   JS truthiness for the usual timer-id-of-0 reason."
  [self]
  (when-not (or (zero? (.-length (retry-of self))) (live-armed? self))
    (unchecked-set self "_liveTimer"
                   (js/setTimeout (fn []
                                    (unchecked-set self "_liveTimer" nil)
                                    ;; the ONE drain nobody awaits, so it is the
                                    ;; one that needs a voice of its own
                                    (fx/tolerate "mailbox retry drain failed" (drain-retry self))
                                    js/undefined)
                                  1500)))
  js/undefined)

(defn- note-gap!
  "A page whose oldest seq is NEWER than our cursor means the node's ring buffer
   dropped messages this device never fetched. Compared with `seq-lt`: a mailbox
   seq is an opaque server string and numeric-parsing it overflows a double."
  [state oldest-seq]
  (when-some [cursor (w/non-empty (:cursor @state))]
    (when (and (some? oldest-seq) (seq-lt cursor oldest-seq))
      (vswap! state assoc :gap? true)))
  js/undefined)

(defn- advance-cursor!
  "The cursor advances only past a FULLY processed page. Re-processing a page is
   idempotent; skipping one loses a message forever."
  [self state msgs]
  (when (pos? (.-length msgs))
    (let [last-seq (unchecked-get (aget msgs (dec (.-length msgs))) "seq")]
      (vswap! state assoc :cursor last-seq)
      (save-str (key-of self "cursor") last-seq)))
  js/undefined)

(defn- process-page!
  "Every message of one page, in arrival order, each failure logged and stepped
   over: one malformed blob must not abandon the rest of the page."
  [self msgs collected]
  (letfn [(step [i]
            (if (>= i (.-length msgs))
              (p/resolved nil)
              (-> (fx/tolerate "mailbox message dropped"
                               (j/later #(process-message self (aget msgs i) collected)))
                  (.then (fn [_] (step (inc i)))))))]
    (step 0)))

(defn- fetch-pages!
  "Phase one: walk the cursor forward, storing every block as it arrives, up to
   MAX-REPLAY-PAGES pages. A transient failure just stops the walk — the next
   kick retries, and the joins below still run on what was collected."
  [self state collected]
  (letfn [(page [n]
            (if (>= n MAX-REPLAY-PAGES)
              (p/resolved nil)
              (p/let [r (.replay ^js (client-of self) (:cursor @state) 100)]
                (let [{:keys [ok? messages more? oldest-seq]} (replay-result r)]
                  ;; transient — the next kick retries; the joins still run
                  (when ok?
                    (note-gap! state oldest-seq)
                    (p/let [_ (process-page! self messages collected)
                            _ (advance-cursor! self state messages)]
                      (when more? (page (inc n)))))))))]
    (page 0)))

(defn- finish-replay!
  "The synchronous tail of a replay, in one tick: tell the user once that the
   ring buffer ate something, then kick the deposit queue (server dedup makes a
   double deposit harmless)."
  [self gap?]
  (when (and (or gap? (pos? (.-length (retry-of self))))
             (not (j/truthy? (unchecked-get self "_noticed"))))
    (unchecked-set self "_noticed" true)
    (on-notice! self "some offline messages expired before this device could fetch them"))
  (flush! self)
  js/undefined)

(defn- replay-all!
  "PROTOCOL.md §9's two phases, in order and never overlapped: fetch and STORE
   every page's blocks first, THEN join what was collected, in seq order —
   which is arrival order, because pages are ordered. That is what stops a join
   from ever waiting on a bitswap timeout.

   `collected` is a JS array because `_processMessage`'s test seam pushes into
   one; `state` is a volatile map so the cursor and the gap flag travel together
   instead of as two closed-over cells."
  [self]
  (let [state (volatile! {:cursor (load-str (key-of self "cursor")) :gap? false})
        collected (array)]
    (p/let [_ (fetch-pages! self state collected)
            joined (join-phase self collected)
            _ (when (pos? joined) (on-replayed! self joined))
            _ (drain-retry self)]
      (finish-replay! self (:gap? @state)))))

(defn- do-replay
  "Replay, throttled to once per REPLAY-THROTTLE-MS unless forced, and never
   re-entrant. `force` is what boot uses — and it is load-bearing: `start`'s
   first replay must not be swallowed by a throttle, or offline delivery is
   silently off for the whole session."
  [self force]
  (if (or (j/truthy? (unchecked-get self "_replaying"))
          (and (not (j/truthy? force))
               (< (- (js/Date.now) (unchecked-get self "_lastReplay")) REPLAY-THROTTLE-MS)))
    (p/resolved nil)
    (do
      (unchecked-set self "_replaying" true)
      (unchecked-set self "_lastReplay" (js/Date.now))
      (-> (replay-all! self)
          ;; `_replaying` is cleared in the `.then` below, and that is why this
          ;; catch exists at all: v24 shipped with both of them detached from
          ;; the chain, so the flag latched replay off permanently
          (.catch (fn [err] (js/console.error "mailbox replay failed" err) nil))
          (.then (fn [_] (unchecked-set self "_replaying" false) js/undefined))))))

(defn- entry-hash
  "A sealed entry's own hash — its address in the log, and its key on the retry
   list."
  [entry]
  (unchecked-get entry "hash"))

(defn- collect-entry!
  "The crash-safety tail, in ONE tick: an entry joins the retry list from the
   moment its block is STORED, before the cursor advances past its page and
   long before any join is attempted. A crash in between must not strand it
   behind an advanced cursor — that is a message lost forever.
   test/vectors/mailbox-effects.json asserts
   putEntryBlock < add-retry < cursor advance < ingestEntry."
  [self entry collected]
  (add-retry self (entry-hash entry))
  (.push collected entry)
  js/undefined)

(defn- store-entry!
  "A dag-cbor block that decodes as a room entry."
  [self block collected]
  (let [^js log (log-of self)]
    (p/let [entry (.decodeSealedEntry log block)
            ;; dag-cbor garbage is not a room entry, and is not looked up
            known (when (j/truthy? entry) (.hasEntry log (entry-hash entry)))]
      (when (and (j/truthy? entry) (not (j/truthy? known)))
        (p/let [_ (.putEntryBlock log (entry-hash entry) block)]
          (collect-entry! self entry collected))))))

(defn- store-block!
  "Store one verified block under the CID it arrived with.

   NEVER TRUST THE SENDER'S CID (PROTOCOL.md §9). The block bytes are re-hashed
   and compared against the multihash digest before anything touches the
   blockstore, and a multihash that is not sha2-256 is refused outright rather
   than verified under whatever it names — otherwise a peer could file its bytes
   under someone else's address."
  [self cid block sealed? collected]
  (let [mh (unchecked-get cid "multihash")]
    (when (identical? (unchecked-get mh "code") SHA2-256)
      (p/let [digest (js/crypto.subtle.digest "SHA-256" block)]
        (when (bytes-eq (js/Uint8Array. digest) (unchecked-get mh "digest"))
          (cond
            sealed? (.putIdentityBlock ^js (log-of self) cid block)
            (identical? (unchecked-get cid "code") DAG-CBOR) (store-entry! self block collected)
            ;; a unixfs image block: a raw leaf or a dag-pb root
            :else (.putRawBlock ^js (log-of self) cid block)))))))

(defn- classify-blob
  "A replayed blob's envelope: `{:sealed? …, :body …}`, or nil when the prefix
   byte is neither 0x01 nor 0x02. Pure, and the ONLY place the prefix bytes are
   read — an unknown prefix is dropped here, before any key is touched."
  [raw]
  (when (>= (.-length raw) 2)
    (let [tag (aget raw 0)]
      (cond
        (identical? tag SEALED-BLOCK) {:sealed? true :body (.subarray raw 1)}
        (identical? tag RAW-BLOCK) {:sealed? false :body (.subarray raw 1)}
        :else nil))))

(defn- decode-cid
  "`CID.decodeFirst` over block bytes, or nil. These bytes are
   attacker-controlled, so a decode failure is a drop, never a throw."
  [body]
  (try (.decodeFirst CID body) (catch :default _ nil)))

(defn- process-message
  "Classify, hash-verify and store one replayed blob; collect entry candidates.

   Every rejection below returns nil and says nothing (PROTOCOL.md §3): a blob
   whose seal will not open is indistinguishable from one that was never meant
   for us, and complaining about it would be both a side channel and a flood."
  [self m collected]
  (j/later
   (fn []
     (when-some [env (classify-blob (rc/from-b64 (unchecked-get m "payload")))]
       (let [sealed? (:sealed? env)]
         (p/let [body (if sealed?
                        (rc/at-rest-open (unchecked-get self "_cipher") (:body env))
                        (:body env))]
           ;; nil body = not ours, or the AEAD tag did not verify
           (when (some? body)
             (when-some [pair (decode-cid body)]
               (store-block! self (aget pair 0) (aget pair 1) sealed? collected)))))))))

(defn- after-ingest!
  "The synchronous verdict on one ingest. Anything that is not `deferred` is
   settled — joined, duplicate or rejected — so it leaves the retry list."
  [self joined entry result]
  (when (identical? "joined" result) (vswap! joined inc))
  (when-not (identical? "deferred" result) (remove-retry self (entry-hash entry)))
  js/undefined)

(defn- project-meta!
  "A joined or deferred entry's display projection. A deferred one is surfaced
   display-only (ChatStore dedups by hash when the real join lands later), and
   an image op's root is remembered so it can be pinned."
  [self img-roots result meta-obj]
  (when (identical? "deferred" result)
    (on-display-only! self meta-obj))
  (let [op (unchecked-get meta-obj "op")]
    (when (and (identical? "img" (unchecked-get op "t"))
               (identical? "string" (js* "typeof ~{}" (unchecked-get op "cid"))))
      (.add img-roots (unchecked-get op "cid"))))
  js/undefined)

(defn- pin-roots!
  "No online source may exist for a mailbox image, so its root is pinned or the
   next gc eats it. Best effort, one root at a time."
  [self roots]
  (j/each-in-order!
   (array-seq roots)
   (fn [root]
     (fx/drop-silently (j/later #(.pinImageIfLocal ^js (log-of self) root))))))

(defn- join-phase
  "Phase two: join every collected entry IN ORDER, project the unjoinable
   display-only, emit once, then pin the image roots the joins ingested.

   Sequential by contract, not by habit: a later entry's ancestors may be an
   earlier entry, so this may never become `p/all`."
  [self entries]
  (let [^js log (unchecked-get self "_log")
        joined (volatile! 0)
        img-roots (js/Set.)]
    (letfn [(step [i]
              (if (>= i (.-length ^js entries))
                (p/resolved nil)
                (let [entry (aget entries i)]
                  (p/let [result (.ingestEntry log entry)
                          _ (after-ingest! self joined entry result)
                          meta-obj (when (or (identical? "joined" result) (identical? "deferred" result))
                                     (.entryMeta log entry))]
                    (when (some? meta-obj) (project-meta! self img-roots result meta-obj))
                    (step (inc i))))))]
      (p/let [_ (step 0)
              _ (when (pos? @joined) (.emitUnseen log))
              _ (pin-roots! self (js/Array.from img-roots))]
        @joined))))

(defn- recheck!
  "One retry-list hash: joined by live replication since (forget it), still
   stored (queue it for another join attempt), or gone to gc (forget it)."
  [self log entries hash]
  (p/let [known (.hasEntry ^js log hash)]
    (if (j/truthy? known)
      (remove-retry self hash)
      (p/let [bytes (.sealedEntryBytes ^js log hash)
              entry (when (j/truthy? bytes) (.decodeSealedEntry ^js log bytes))]
        (if (j/truthy? entry)
          (do (.push entries entry) js/undefined)
          (remove-retry self hash))))))

(defn- report-joined! [self joined]
  (when (pos? joined) (on-replayed! self joined))
  js/undefined)

(defn- drain-retry
  "Re-attempt stored-but-unjoined entries; their blocks are already local, so
   this costs no network. Walks a SNAPSHOT because `recheck!` mutates the live
   list underneath it."
  [self]
  (if (zero? (.-length (retry-of self)))
    (p/resolved nil)
    (let [log (log-of self)
          snapshot (.slice (retry-of self))
          entries (array)]
      (p/let [_ (j/each-in-order! (array-seq snapshot) #(recheck! self log entries %))
              joined (join-phase self entries)]
        (report-joined! self joined)))))

(defn- add-retry [self hash]
  (when-not (.includes (retry-of self) hash)
    (.push (retry-of self) hash)
    (save-json (key-of self "retry") (retry-of self)))
  js/undefined)

(defn- remove-retry [self hash]
  (let [i (.indexOf (retry-of self) hash)]
    (when-not (identical? -1 i)
      (.splice (retry-of self) i 1)
      (save-json (key-of self "retry") (retry-of self))))
  js/undefined)

;; ---- class surface ---------------------------------------------------------
;;
;; Frozen by test/vectors/facade.json: the names, their ORDER and their arities.
;; The twelve `_` methods are the seam test/helpers/mailbox-effects.mjs drives —
;; the only way to say "and then the retry timer fired 2 seconds later".

(let [proto (.-prototype MailboxSync)]
  (js/Object.defineProperty
   proto "state"
   (j/ordered "get" (fn [] (this-as self (state self))) "configurable" true))
  (unchecked-set proto "start" (fn [] (this-as self (start self))))
  (unchecked-set proto "stop" (fn [] (this-as self (stop self))))
  (unchecked-set proto "enqueueEntry" (fn [hash] (this-as self (enqueue-entry self hash))))
  (unchecked-set proto "enqueueImage" (fn [root] (this-as self (enqueue-image self root))))
  (unchecked-set proto "onLiveEntry" (fn [] (this-as self (on-live-entry self))))
  ;; the "private" surface the tests drive (peer-kit's `_name` idiom)
  (unchecked-set proto "_buildPayload" (fn [item] (this-as self (build-payload self item))))
  (unchecked-set proto "_processMessage" (fn [m collected] (this-as self (process-message self m collected))))
  (unchecked-set proto "_enqueueIdentityOnce" (fn [] (this-as self (enqueue-identity-once self))))
  (unchecked-set proto "_push" (fn [item] (this-as self (push-item self item))))
  (unchecked-set proto "_dropHead" (fn [item] (this-as self (drop-head self item))))
  (unchecked-set proto "_markSent" (fn [k] (this-as self (mark-sent self k))))
  (unchecked-set proto "_flush" (fn [] (this-as self (flush! self))))
  (unchecked-set proto "_replay" (fn [force] (this-as self (do-replay self force))))
  (unchecked-set proto "_drainRetry" (fn [] (this-as self (drain-retry self))))
  (unchecked-set proto "_joinPhase" (fn [entries] (this-as self (join-phase self entries)))))

static mirror of HEAD · about · clone: git clone https://git.ardegazu.ro/chat.git