rooms-kit / src / ardegazu / rooms / lib / log.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
;; ported-from: src/lib/log.ts @ v1.0.0
;;
;; The room's replicated, encrypted, append-only log (OrbitDB events DB on Helia,
;; PROTOCOL.md v2). This is the only namespace that touches @orbitdb/core.
;;
;; Address determinism: every member derives the same db name from the room
;; secret and passes an identical access controller, so the manifest CID — and
;; therefore the address — is equal for all of them with no exchange needed.
;; Access control IS the encryption: write:["*"] plus decrypt-or-drop.
;;
;; History model: replication replaces v1's one-shot hist-req pull. Entries
;; persist in the device blockstore, so reopening a room offline replays it from
;; disk, and any member's device serves late joiners.
;;
;; SHAPE. The file reads top-down in dependency order: the plain-value helpers
;; first, then the update chain, then opening, then the block-level mailbox
;; surface. That ordering is the point — every `defn` here is defined before its
;; first use, so this namespace carries NO forward declarations. The `(declare)`
;; it used to open with hid the fact that the projection, the chain and the
;; ingest path do not actually depend on each other in a circle.
;;
;; INTEROP NOTES for the port (this file is almost entirely interop):
;;  - `multiformats` is imported but deliberately NOT a declared dependency. It
;;    reaches this build transitively through helia/@orbitdb/core, which
;;    guarantees ONE copy: CID instances cross into helia and @orbitdb/core, and
;;    two copies of multiformats would break `instanceof CID` inside them. v1.0.0
;;    relied on the same transitive resolution.
;;  - every async generator (the oplog iterator, unixfs `cat`, `pins.add`/`rm`)
;;    is driven through ardegazu.rooms.js's async-iteration helpers; CLJS has no
;;    `for await`.
;;  - the promise chains go through ardegazu.rooms.js's `later` / `attempt` /
;;    `each-in-order!` / `every-in-order!`. `attempt` IS the decrypt-or-drop rule
;;    (§3) written once: any failure, synchronous throw included, resolves to the
;;    fallback. The two fold helpers are sequential by contract — the order these
;;    side effects run in is protocol, not implementation detail.
;;  - `entry?.hash` and `entry?.clock?.time ?? 0` are JS truthiness/nullish, not
;;    Clojure's — see ardegazu.rooms.js.
(ns ardegazu.rooms.lib.log
  (:require ["@helia/block-brokers" :refer (bitswap)]
            ["@helia/routers" :refer (libp2pRouting)]
            ["@helia/unixfs" :refer (unixfs)]
            ["helia" :refer (createHelia)]
            ["@ipld/dag-pb" :as dag-pb]
            ["ipfs-unixfs-importer/chunker" :refer (fixedSize)]
            ["multiformats/cid" :refer (CID)]
            ["multiformats/bases/base58" :refer (base58btc)]
            ["@orbitdb/core" :as orbit]
            [ardegazu.rooms.js :as j]
            [ardegazu.rooms.lib.access :as access]
            [ardegazu.rooms.lib.crypto :as rc]
            [ardegazu.rooms.lib.encryption :as encryption]
            [ardegazu.rooms.lib.orbit-identity :as oi]
            [shadow.cljs.modern :refer (defclass js-await)]))

(def ^:private IMG-MAX-BYTES (+ (* 2 1024 1024) 1024)) ; re-encode cap + envelope slack
(def ^:private DAG-PB 0x70) ; dag-pb multicodec (unixfs internal nodes; leaves are raw)
(def ^:private TAIL-AMOUNT 500) ; entries a straggler sweep reads back
(def ^:private MAX-ANCESTRY 5000) ; runaway chain — let live replication carry it

;; IngestResult — the outcome of feeding one mailbox-delivered entry into the log:
;;   "joined"    verified and in the log — will surface via emitUnseen()
;;   "duplicate" already in the log
;;   "deferred"  an ancestor or identity block is not local yet — retry later
;;   "rejected"  failed verification (signature/access/log id) — drop for good

(defn open-helia
  "One Helia node per session, shared across rooms. The stores are CONSTRUCTED
   AND OPENED BY THE CALLER and handed in already open — this namespace is shared
   with the browser builds, where they are IndexedDB-backed, while a Node host's
   are filesystem-backed; neither package may be required here or the browser
   bundle would pull Node's fs in. Whichever flavour, two apps must not share a
   blockstore: bitswap serves every local block to connected peers, and GC/quota
   would mix across apps. Node callers: ardegazu.rooms.node.stores/open-fs-helia
   is this function behind the historical store-PATHS signature."
  [libp2p blockstore datastore]
  ;; libp2p is created and owned by Net; Helia must not re-create or stop
  ;; it. Bitswap-only, member↔member: Helia's DEFAULTS also race public
  ;; HTTP gateways (trustless-gateway.link, 4everland.io) for every block —
  ;; our blocks are ciphertext held only by members, so those requests
  ;; would leak CIDs + client IPs to third parties for zero benefit. Never
  ;; add them back.
  (createHelia (j/ordered "libp2p" libp2p
                          "blockstore" blockstore
                          "datastore" datastore
                          "blockBrokers" #js [(bitswap)]
                          "routers" #js [(libp2pRouting libp2p)]
                          "start" true)))

;; A LogEntryMeta is {hash, from, clock, op}:
;;   from  — verified author id: base64url Ed25519 pub for sueta identities, a
;;           stable per-device key for OrbitDB's default "publickey" identities,
;;           "" when the entry's signing key doesn't match its referenced
;;           identity (see author-of).
;;   clock — Lamport time, used only for folding reaction ops deterministically.

(defclass RoomLog
  (constructor [this db orbitdb helia img-cipher encryption]
    (unchecked-set this "_db" db)
    (unchecked-set this "_orbitdb" orbitdb)
    (unchecked-set this "_helia" helia)
    (unchecked-set this "_fs" (unixfs helia))
    (unchecked-set this "_imgCipher" img-cipher)
    (unchecked-set this "_encryption" encryption)
    (unchecked-set this "_authorCache" (js/Map.)) ; identity hash → record
    (unchecked-set this "_seen" (js/Set.))
    (unchecked-set this "_updateChain" (js/Promise.resolve nil))
    ;; Live + replicated entries, deduped by hash. Wired inside open() BEFORE
    ;; replication starts — entries can arrive the moment sync.start() runs, and
    ;; emit marks hashes as seen, so a late-attached handler would lose them
    ;; permanently. Pass both callbacks to RoomLog.open().
    (unchecked-set this "onEntry" (fn [_] js/undefined))
    (unchecked-set this "onError" (fn [_] js/undefined))
    ;; Mailbox hooks (wired by the app when the feature is on; default no-ops).
    (unchecked-set this "onAppended" (fn [_] js/undefined))
    (unchecked-set this "onImageStored" (fn [_] js/undefined))
    ;; When set, putImage chunks at this size so every unixfs block fits one
    ;; mailbox frame. Only affects NEW images; the CID travels in the entry.
    (unchecked-set this "imageChunkBytes" nil)))

;; ---- the objects hanging off `self` ----------------------------------------
;;
;; dev/docs/CLJS.md's externs pitfall: a `^js` hint on a MACRO form
;; (unchecked-get) is LOST, so `(.joinEntry (unchecked-get … "log") e)` can never
;; be inferred. Every interop call below therefore binds its receiver to a
;; `^js`-tagged local first. These five accessors are what makes that one line
;; instead of three at each of the ~20 call sites.

(defn- oplog [self] (unchecked-get (unchecked-get self "_db") "log"))
(defn- entry-storage
  "The oplog's IPFSBlockStorage — sealed entry blocks, pinned by OrbitDB."
  [self]
  (unchecked-get (oplog self) "storage"))
(defn- blockstore [self] (unchecked-get (unchecked-get self "_helia") "blockstore"))
(defn- pin-store [self] (unchecked-get (unchecked-get self "_helia") "pins"))
(defn- log-iterator
  "The oplog's newest→oldest async iterator over the last `amount` entries."
  [self amount]
  (let [^js l (oplog self)]
    (.iterator l (j/ordered "amount" amount))))

;; ---- plain-value helpers ---------------------------------------------------

(defn- op-of
  "The events DB wraps each value as {op:'ADD', key:null, value}; older/raw
   payloads are the op itself. nil when the payload is not a typed op."
  [entry]
  (let [wrapper (unchecked-get entry "payload")
        op (if (and (j/truthy? wrapper) (identical? "ADD" (unchecked-get wrapper "op")))
             (unchecked-get wrapper "value")
             wrapper)]
    (when (and (j/truthy? op)
               (j/js-object? op)
               (identical? "string" (js* "typeof ~{}" (unchecked-get op "t"))))
      op)))

(defn- clock-time
  "Lamport time of an entry, 0 when it carries none. TS `entry.clock?.time ?? 0`
   — optional chaining then NULLISH, neither of which is Clojure's nil-punning."
  [entry]
  (let [c (unchecked-get entry "clock")]
    (j/nn (when (some? c) (unchecked-get c "time")) 0)))

(defn- meta-row
  "The LogEntryMeta the app sees. Built in key order by policy (dev/docs/CLJS.md)
   and shared by `emit` and `entryMeta`, which must produce the same shape — the
   whole point of entryMeta is that a later real join emits an identical row."
  [entry from op]
  (j/ordered "hash" (unchecked-get entry "hash")
             "from" from
             "clock" (clock-time entry)
             "op" op))

(defn- field-array
  "A JS array field of `o`, or an empty array when it is absent — `next` and
   `refs` are both optional on an entry, and both are iterated the same way."
  [o k]
  (let [v (when (j/truthy? o) (unchecked-get o k))]
    (if (j/truthy? v) v (array))))

(defn- all-seen?
  "Has every hash in `hashes` already been emitted?"
  [^js seen hashes]
  (not (.some ^js hashes (fn [h] (not (.has seen h))))))

;; ---- author binding --------------------------------------------------------

(defn- trusted-author
  "The {id, publicKey} pair an identity record binds, or the empty pair. \"sueta\"
   identities yield the b64url Ed25519 pub (fingerprint/TOFU-able); OrbitDB's
   default \"publickey\" identities yield a stable per-device id so identity-less
   users still get working \"mine\" styling and reactions. Any other provider
   type is not trusted for authorship."
  [ident]
  (if (and (j/truthy? ident)
           (or (identical? "sueta" (unchecked-get ident "type"))
               (identical? "publickey" (unchecked-get ident "type"))))
    (j/ordered "id" (unchecked-get ident "id")
               "publicKey" (unchecked-get ident "publicKey"))
    (j/ordered "id" "" "publicKey" "")))

(defn- resolve-author
  "The identity record `ref` names, memoised for the session. FAILURES ARE CACHED
   TOO — an identity that cannot be resolved must not be re-fetched once per
   entry that references it."
  [self ref]
  (let [^js cache (unchecked-get self "_authorCache")
        cached (.get cache ref)]
    (if-not (identical? cached js/undefined)
      (js/Promise.resolve cached)
      (-> (j/attempt (fn []
                       (js-await [ident (let [^js ids (unchecked-get (unchecked-get self "_orbitdb") "identities")]
                                          (.getIdentity ids ref))]
                         (trusted-author ident)))
                     (j/ordered "id" "" "publicKey" ""))
          (.then (fn [rec] (.set cache ref rec) rec))))))

(defn- author-of
  "Durable author id of an entry, defense-in-depth on top of the hardened access
   controller: the binding only counts when the entry's signing key IS the
   referenced identity's key (lib/access closes this at replication; here it also
   covers entries accepted by older clients)."
  [self entry]
  (let [ref (unchecked-get entry "identity")
        signing-key (unchecked-get entry "key")]
    (if (or (not (j/truthy? ref)) (not (j/truthy? signing-key)))
      (js/Promise.resolve "")
      (.then (resolve-author self ref)
             (fn [rec]
               (if (and (j/truthy? (unchecked-get rec "id"))
                        (identical? (unchecked-get rec "publicKey") signing-key))
                 (unchecked-get rec "id")
                 ""))))))

;; ---- projection and the update chain ---------------------------------------

(defn- emit
  "Project one entry to the app, at most once per hash. `_seen` is marked BEFORE
   the author lookup awaits, so a second call for the same hash arriving inside
   that window is still deduped."
  [self entry]
  (let [^js seen (unchecked-get self "_seen")
        hash (when (j/truthy? entry) (unchecked-get entry "hash"))]
    (if (or (not (j/truthy? hash)) (.has seen hash))
      (js/Promise.resolve nil)
      (do
        (.add seen hash)
        (let [op (op-of entry)]
          (if (nil? op)
            (js/Promise.resolve nil)
            (js-await [from (author-of self entry)]
              (do ((unchecked-get self "onEntry") (meta-row entry from op))
                  nil))))))))

(defn- sweep-unseen
  "Emit every entry in the log the app has not seen yet. The oplog iterator
   yields newest→oldest, so the batch is REVERSED before emitting: the app must
   build its state oldest-first."
  [self]
  (let [unseen (array)
        ^js seen (unchecked-get self "_seen")]
    (-> (j/for-each! (log-iterator self TAIL-AMOUNT)
                     (fn [e]
                       (when-not (.has seen (unchecked-get e "hash")) (.push unseen e))
                       nil))
        (.catch (fn [err] ((unchecked-get self "onError") err) nil))
        (.then (fn [_]
                 (.reverse unseen)
                 (j/each-in-order! (array-seq unseen) (fn [e] (emit self e))))))))

(defn- chain!
  "Append `f` to the serialized update chain and return the NEW chain. Live
   updates, straggler sweeps and mailbox replay all queue here, which is what
   stops them interleaving."
  [self f]
  (let [next-chain (.then (unchecked-get self "_updateChain") f)]
    (unchecked-set self "_updateChain" next-chain)
    next-chain))

(defn- on-update
  "'update' fires once per new HEAD — a replicated branch join delivers many
   entries under one event. If the head links to entries we haven't emitted,
   sweep the log for stragglers (they're all in local storage by the time the
   event fires).

   TWO THINGS HERE ARE CONTRACT, not style. The `_seen` test is read at RUN time,
   inside the queued callback, never hoisted: two updates queued back to back
   where the second links to the first must produce ONE sweep, because by the
   time the second callback runs the first has already marked its hash. And this
   is fire-and-forget — it returns undefined, not the chain. Its sibling
   `emit-unseen` returns the chain, and lib/mailbox depends on the difference."
  [self entry]
  (chain! self
          (fn [_]
            (if (all-seen? (unchecked-get self "_seen") (field-array entry "next"))
              (emit self entry)
              (.then (sweep-unseen self) (fn [_] (emit self entry))))))
  js/undefined)

(defn emit-unseen
  "Surface entries added via ingestEntry — a direct joinEntry doesn't fire the
   'update' event, so mailbox replay calls this once after its joins. Returns THE
   CHAIN, and the identical object `_updateChain` now holds: lib/mailbox's
   join-phase awaits this before reporting how many entries joined, so anything
   else here would let it report and flush before the entries had surfaced."
  [self]
  (chain! self (fn [_] (sweep-unseen self))))

;; ---- opening ---------------------------------------------------------------

(defn- log-directory
  "The room's state directory. TS `hooks?.directory ?? \"sueta\"` — NULLISH, so an
   explicitly empty directory string is honoured rather than replaced. The
   directory embeds the salt-derived roomId, so it is already unique per app;
   callers pass a base under their state dir (Node port — the browser original
   used a relative IndexedDB-backed path)."
  [room-crypto hooks]
  (let [base (j/nn (when (j/truthy? hooks) (unchecked-get hooks "directory")) "sueta")]
    (str base "/" (unchecked-get room-crypto "roomId"))))

(defn- open-identities
  "OrbitDB's identities store plus this device's sueta identity, or nil when the
   caller passed none — OrbitDB then falls back to its own \"publickey\" provider."
  [helia identity directory]
  (if-not (j/truthy? identity)
    (js/Promise.resolve nil)
    (do
      (oi/register-sueta-provider)
      (js-await [keystore (orbit/KeyStore (j/ordered "path" (str directory "/keystore")))]
        (js-await [identities (orbit/Identities (j/ordered "ipfs" helia "keystore" keystore))]
          (js-await [orbit-identity (.createIdentity ^js identities
                                                     (j/ordered "provider"
                                                                (oi/SuetaIdentityProvider
                                                                 (j/ordered "identity" identity))))]
            (j/ordered "identities" identities "identity" orbit-identity)))))))

(defn- orbit-options
  "createOrbitDB's options, with the identity pair spliced in when there is one."
  [helia directory ids]
  (let [opts (j/ordered "ipfs" helia "directory" directory)]
    (when (j/truthy? ids)
      (unchecked-set opts "identities" (unchecked-get ids "identities"))
      (unchecked-set opts "identity" (unchecked-get ids "identity")))
    opts))

(defn- open-db
  "The room's events DB.

   sync:false → the caller attaches the error listener BEFORE any peer exchange
   can throw (a non-member's undecryptable garbage must never crash us).

   Hardened wrapper: same manifest address as the vanilla IPFS controller (only
   {type:\"ipfs\", write} is hashed), but forged-authorship entries — signed by
   one key while referencing another key's identity — are rejected at replication
   time. See lib/access. write:[\"*\"] is address-determining: changing it forks
   the DB."
  [orbitdb db-name encryption]
  (.open ^js orbitdb db-name
         (j/ordered "type" "events"
                    "AccessController" (access/HardenedIPFSAccessController
                                        (j/ordered "write" #js ["*"]))
                    "encryption" encryption
                    "sync" false)))

(defn- open-orbit
  "Everything between a Helia node and an open events DB, in the order the
   original ladder ran it: identities, OrbitDB, the derived db name, the log
   encryption, the DB itself."
  [helia room-crypto identity directory]
  (js-await [ids (open-identities helia identity directory)]
    (js-await [orbitdb (orbit/createOrbitDB (orbit-options helia directory ids))]
      (js-await [db-name (rc/db-name room-crypto)]
        (js-await [encryption (encryption/make-log-encryption room-crypto)]
          (js-await [db (open-db orbitdb db-name encryption)]
            (j/ordered "orbitdb" orbitdb "db" db "encryption" encryption)))))))

(defn- wire-hooks!
  "Attach the app's callbacks BEFORE replication starts: entries can arrive the
   moment sync.start() runs, and emit marks hashes as seen, so a late-attached
   handler would lose them permanently."
  [log db hooks]
  (let [^js events (unchecked-get db "events")
        on-entry (when (j/truthy? hooks) (unchecked-get hooks "onEntry"))
        on-error (when (j/truthy? hooks) (unchecked-get hooks "onError"))]
    (when (j/truthy? on-entry) (unchecked-set log "onEntry" on-entry))
    (when (j/truthy? on-error) (unchecked-set log "onError" on-error))
    (.on events "error" (fn [err] ((unchecked-get log "onError") err)))
    (.on events "update" (fn [entry] (on-update log entry))))
  log)

(defn open [helia room-crypto identity hooks]
  (j/later
   (fn []
     (js-await [parts (open-orbit helia room-crypto identity (log-directory room-crypto hooks))]
       (js-await [img-cipher (rc/img-cipher room-crypto)]
         (let [db (unchecked-get parts "db")
               log (wire-hooks! (RoomLog. db (unchecked-get parts "orbitdb") helia
                                          img-cipher (unchecked-get parts "encryption"))
                                db hooks)]
           (js-await [_ (.start ^js (unchecked-get db "sync"))]
             log)))))))

;; ---- appending / replay ----------------------------------------------------

(defn append [self op]
  (js-await [hash (.add ^js (unchecked-get self "_db") op)]
    (do
      ;; mailbox enqueue must never fail an append
      (try ((unchecked-get self "onAppended") hash) (catch :default _ nil))
      hash)))

(defn load-tail
  "Replay the newest `amount` entries from local storage (offline-capable);
   resolves to how many were read."
  [self amount]
  (let [entries (array)]
    (-> (j/for-each! (log-iterator self amount) (fn [e] (.push entries e) nil))
        (.catch (fn [err] ((unchecked-get self "onError") err) nil))
        (.then (fn [_]
                 ;; iterator yields newest→oldest; emit oldest-first so the
                 ;; consumer builds state in order
                 (.reverse entries)
                 (j/each-in-order! (array-seq entries) (fn [e] (emit self e)))))
        (.then (fn [_] (.-length entries))))))

;; ---- images: encrypted unixfs blobs, replicated via bitswap ---------------

(defn- concat-bytes
  "Join `parts` (Uint8Arrays, `total` bytes between them) into one Uint8Array."
  [parts total]
  (let [out (js/Uint8Array. total)]
    (loop [i 0 off 0]
      (if (< i (.-length ^js parts))
        (let [p (aget parts i)]
          (.set out p off)
          (recur (inc i) (+ off (.-length ^js p))))
        out))))

(defn put-image [self data]
  (js-await [sealed (rc/at-rest-seal (unchecked-get self "_imgCipher") data)]
    (js-await [cid (let [^js fs (unchecked-get self "_fs")
                         chunk (unchecked-get self "imageChunkBytes")]
                     (.addBytes fs sealed
                                (if (j/truthy? chunk)
                                  (j/ordered "chunker" (fixedSize (j/ordered "chunkSize" chunk)))
                                  (j/ordered))))]
      (-> (j/attempt (fn [] (let [^js ps (pin-store self)] (j/drain! (.add ps cid))))
                     nil) ; already pinned
          (.then (fn [_]
                   (try ((unchecked-get self "onImageStored") (.toString cid))
                        (catch :default _ nil))
                   (j/ordered "cid" (.toString cid) "bytes" (.-length ^js sealed))))))))

(defn get-image
  "Fetch + decrypt an image blob (local blockstore or bitswap from members)."
  [self cid-str signal]
  (j/attempt
   (fn []
     (let [cid (.parse CID cid-str)
           parts (array)
           total (volatile! 0)]
       (js-await [_ (j/for-each! (let [^js fs (unchecked-get self "_fs")]
                                   (.cat fs cid (j/ordered "signal" signal)))
                                 (fn [chunk]
                                   (vswap! total + (.-length ^js chunk))
                                   ;; liar: bigger than the protocol cap
                                   (when (> @total IMG-MAX-BYTES)
                                     (throw (js/Error. "img too big")))
                                   (.push parts chunk)
                                   nil))]
         (rc/at-rest-open (unchecked-get self "_imgCipher") (concat-bytes parts @total)))))
   nil))

(defn prune-images
  "Unpin images that fell out of the newest-N window; blocks GC'd next sweep."
  [self keep dropped]
  (-> (j/each-in-order!
       (array-seq dropped)
       (fn [cid-str]
         (if (.has ^js keep cid-str)
           nil
           (j/attempt (fn [] (let [^js ps (pin-store self)]
                               (j/drain! (.rm ps (.parse CID cid-str)))))
                      nil)))) ; not pinned here — fine
      (.then (fn [_] (.catch (js/Promise.resolve (.gc ^js (unchecked-get self "_helia")))
                             (fn [_] nil))))
      (.then (fn [_] js/undefined))))

;; ---- offline mailbox surface (PROTOCOL.md §9) -----------------------------
;;
;; Export: hand the orchestrator the exact stored bytes (sealed entry blocks,
;; unixfs image blocks, this device's identity record). Import: verify and feed
;; replayed blocks back in with the SAME guarantees as live replication
;; (Entry.decode = decrypt-or-drop; joinEntry = signature + hardened access
;; control), never letting a missing block reach a bitswap timeout.

(defn sealed-entry-bytes
  "Sealed (kEntry) block bytes of a local entry — byte-identical on every replica.
   nil when absent (the oplog's storage resolves undefined, TS coerced with ??)."
  [self hash]
  (j/attempt (fn []
               (js-await [b (let [^js storage (entry-storage self)] (.get storage hash))]
                 (j/nn b nil)))
             nil))

(defn my-identity-block
  "This device's plaintext identity record block (sealed by the caller for
   transit); nil when the record's hash is not a parseable base58btc CID."
  [self]
  (try
    (let [ident (unchecked-get (unchecked-get self "_orbitdb") "identity")]
      (j/ordered "cid" (.parse CID (unchecked-get ident "hash") base58btc)
                 "bytes" (unchecked-get ident "bytes")))
    (catch :default _ nil)))

(defn raw-block
  "A locally stored block, never touching bitswap; nil when absent."
  [self cid]
  (j/attempt (fn []
               (let [^js bs (blockstore self)]
                 (js-await [present (.has bs cid)]
                   (if-not (j/truthy? present) nil (.get bs cid)))))
             nil))

(defn- dag-links
  "The child CIDs a dag-pb block links to, or nil when the bytes do not decode."
  [bytes]
  (try (.map ^js (unchecked-get (dag-pb/decode bytes) "Links") (fn [l] (unchecked-get l "Hash")))
       (catch :default _ nil)))

(defn image-dag-blocks
  "All blocks of a locally stored image DAG, children before parents (so a
   replayer that processes in order always has leaves before the root pin). []
   when any block is missing or the DAG is deeper than unixfs produces at our
   chunk sizes (root + raw leaves)."
  [self root-cid]
  (let [out (array)]
    (letfn [(keep! [cid bytes]
              (.push out (j/ordered "cid" cid "bytes" bytes))
              true)
            (walk [cid depth]
              (js-await [bytes (raw-block self cid)]
                (cond
                  (not (j/truthy? bytes)) false
                  (not (identical? DAG-PB (unchecked-get cid "code"))) (keep! cid bytes)
                  ;; deeper than root+leaves — refuse, don't ship a partial DAG
                  (>= depth 1) false
                  :else
                  (let [links (dag-links bytes)]
                    (if (nil? links)
                      false
                      (js-await [ok (j/every-in-order! (array-seq links)
                                                       (fn [l] (walk l (inc depth))))]
                        (if-not ok false (keep! cid bytes))))))))]
      (-> (j/later (fn [] (walk (.parse CID root-cid) 0)))
          (.then (fn [ok] (if ok out (array))))
          (.catch (fn [_] (array)))))))

(defn decode-sealed-entry
  "Decode a sealed entry block; nil on garbage (not a member's bytes / tampered)."
  [self bytes]
  (j/attempt
   (fn []
     (let [enc (unchecked-get self "_encryption")]
       (js-await [entry (.decode orbit/Entry bytes
                                 (unchecked-get (unchecked-get enc "replication") "decrypt")
                                 (unchecked-get (unchecked-get enc "data") "decrypt"))]
         (if (and (j/truthy? entry)
                  (j/truthy? (unchecked-get entry "hash"))
                  (j/truthy? (.isEntry orbit/Entry entry)))
           entry
           nil))))
   nil))

(defn put-entry-block
  "Store a sealed entry block (pinned via the oplog's IPFSBlockStorage)."
  [self hash bytes]
  (.put ^js (entry-storage self) hash bytes))

(defn put-raw-block
  "Store an image block; pinning happens per-DAG in pinImageIfLocal."
  [self cid bytes]
  (.put ^js (blockstore self) cid bytes))

(defn put-identity-block
  "Store + pin an identity record block (tiny, link-free — the pin never leaves
   disk)."
  [self cid bytes]
  (js-await [_ (let [^js bs (blockstore self)] (.put bs cid bytes))]
    (-> (j/attempt (fn [] (let [^js ps (pin-store self)] (j/drain! (.add ps cid))))
                   nil) ; already pinned
        (.then (fn [_] js/undefined)))))

(defn has-entry [self hash]
  (j/attempt (fn [] (let [^js l (oplog self)] (.has l hash))) false))

(defn- identity-local
  "Is the entry's identity record resolvable without the network?"
  [self e]
  (let [ref (unchecked-get e "identity")
        ^js cache (unchecked-get self "_authorCache")]
    (cond
      (not (j/truthy? ref)) (js/Promise.resolve false)
      (.has cache ref) (js/Promise.resolve true)
      :else (j/attempt (fn [] (let [^js bs (blockstore self)]
                                (.has bs (.parse CID ref base58btc))))
                       false))))

(defn- can-join-locally
  "joinEntry's own traversal (next ∪ refs, stopping at entries already in the
   log), plus the requirement that every hop's author identity block is already
   local — joinEntry's canAppend would otherwise fetch it with a 30 s bitswap
   timeout."
  [self entry]
  (let [visited (js/Set. #js [(unchecked-get entry "hash")])
        frontier (array)]
    (letfn [(expand [e]
              (js-await [ok (identity-local self e)]
                (if-not (j/truthy? ok)
                  false
                  (do
                    (doseq [h (concat (array-seq (field-array e "next"))
                                      (array-seq (field-array e "refs")))]
                      (when-not (.has visited h)
                        (.add visited h)
                        (.push frontier h)))
                    true))))
            (loop-step []
              (cond
                (identical? 0 (.-length frontier)) (js/Promise.resolve true)
                ;; runaway chain — let live replication carry it
                (> (.-size visited) MAX-ANCESTRY) (js/Promise.resolve false)
                :else
                (let [hash (.pop frontier)]
                  (js-await [known (has-entry self hash)]
                    (if (j/truthy? known)
                      (loop-step) ; joined — its ancestry is in the log
                      (js-await [bytes (j/attempt (fn [] (raw-block self (.parse CID hash base58btc)))
                                                  nil)]
                        (if-not (j/truthy? bytes)
                          false
                          (js-await [e (decode-sealed-entry self bytes)]
                            (if-not (j/truthy? e)
                              false
                              (js-await [ok (expand e)]
                                (if-not ok false (loop-step))))))))))))]
      (js-await [ok (expand entry)]
        (if-not ok false (loop-step))))))

(defn ingest-entry
  "Feed one decoded mailbox entry into the log. Anything not fully resolvable
   locally is \"deferred\" for retry; a throw past the pre-check is a
   verification failure ⇒ \"rejected\"."
  [self entry]
  (js-await [dup (has-entry self (unchecked-get entry "hash"))]
    (if (j/truthy? dup)
      "duplicate"
      (js-await [joinable (can-join-locally self entry)]
        (if-not (j/truthy? joinable)
          "deferred"
          (j/attempt (fn []
                       (js-await [added (let [^js l (oplog self)] (.joinEntry l entry))]
                         (if (j/truthy? added) "joined" "duplicate")))
                     "rejected"))))))

(defn entry-meta
  "Display-only projection of an entry that can't be joined yet (missing
   ancestors). Same shape emit produces; NOT marked seen, so a later real join
   still emits it (consumers dedupe by hash)."
  [self entry]
  (let [op (op-of entry)]
    (if (nil? op)
      (js/Promise.resolve nil)
      (js-await [local (identity-local self entry)]
        (js-await [from (if (j/truthy? local) (author-of self entry) (js/Promise.resolve ""))]
          (meta-row entry from op))))))

(defn pin-image-if-local
  "Pin a mailbox-ingested image root once its whole DAG is local. Receivers never
   pin bitswap-fetched images (they can refetch from members), but a mailbox
   image may have NO online source — without the pin, pruneImages' gc would
   delete in-window blocks. The window's unpin path already covers these
   (pruneImages takes any CID)."
  [self root-cid]
  (js-await [blocks (image-dag-blocks self root-cid)] ; [] unless fully local
    (if (identical? 0 (.-length ^js blocks))
      false
      (-> (j/attempt (fn [] (let [^js ps (pin-store self)]
                              (j/drain! (.add ps (.parse CID root-cid)))))
                     nil) ; already pinned
          (.then (fn [_] true))))))

(defn close [self]
  (js-await [_ (.catch (.close ^js (unchecked-get self "_db")) (fn [_] nil))]
    (.catch (.stop ^js (unchecked-get self "_orbitdb")) (fn [_] nil))))

(defn drop-log
  "Forget-room: drop the log's local data (entries, index, heads)."
  [self]
  (js-await [_ (let [^js db (unchecked-get self "_db")]
                 (.catch (.drop db) (fn [_] nil)))]
    (.catch (.stop ^js (unchecked-get self "_orbitdb")) (fn [_] nil))))

;; ---- class surface ---------------------------------------------------------

(unchecked-set RoomLog "open"
               (fn [helia room-crypto identity hooks] (open helia room-crypto identity hooks)))

(let [proto (.-prototype RoomLog)]
  ;; string-named getters: rename-safe at every optimization level
  (js/Object.defineProperty
   proto "address"
   (j/ordered "get" (fn [] (this-as self (unchecked-get (unchecked-get self "_db") "address")))
              "configurable" true))
  (js/Object.defineProperty
   proto "myAuthorId"
   (j/ordered "get" (fn [] (this-as self (unchecked-get (unchecked-get (unchecked-get self "_orbitdb") "identity") "id")))
              "configurable" true))
  (unchecked-set proto "emitUnseen" (fn [] (this-as self (emit-unseen self))))
  (unchecked-set proto "append" (fn [op] (this-as self (append self op))))
  (unchecked-set proto "loadTail" (fn [amount] (this-as self (load-tail self amount))))
  (unchecked-set proto "putImage" (fn [data] (this-as self (put-image self data))))
  (unchecked-set proto "getImage" (fn [cid signal] (this-as self (get-image self cid signal))))
  (unchecked-set proto "pruneImages" (fn [keep dropped] (this-as self (prune-images self keep dropped))))
  (unchecked-set proto "sealedEntryBytes" (fn [hash] (this-as self (sealed-entry-bytes self hash))))
  (unchecked-set proto "myIdentityBlock" (fn [] (this-as self (my-identity-block self))))
  (unchecked-set proto "rawBlock" (fn [cid] (this-as self (raw-block self cid))))
  (unchecked-set proto "imageDagBlocks" (fn [root] (this-as self (image-dag-blocks self root))))
  (unchecked-set proto "decodeSealedEntry" (fn [bytes] (this-as self (decode-sealed-entry self bytes))))
  (unchecked-set proto "putEntryBlock" (fn [hash bytes] (this-as self (put-entry-block self hash bytes))))
  (unchecked-set proto "putRawBlock" (fn [cid bytes] (this-as self (put-raw-block self cid bytes))))
  (unchecked-set proto "putIdentityBlock" (fn [cid bytes] (this-as self (put-identity-block self cid bytes))))
  (unchecked-set proto "hasEntry" (fn [hash] (this-as self (has-entry self hash))))
  (unchecked-set proto "ingestEntry" (fn [entry] (this-as self (ingest-entry self entry))))
  (unchecked-set proto "entryMeta" (fn [entry] (this-as self (entry-meta self entry))))
  (unchecked-set proto "pinImageIfLocal" (fn [root] (this-as self (pin-image-if-local self root))))
  (unchecked-set proto "close" (fn [] (this-as self (close self))))
  (unchecked-set proto "drop" (fn [] (this-as self (drop-log self))))
  ;; the "private" surface the tests drive (peer-kit's `_name` idiom)
  (unchecked-set proto "_emit" (fn [entry] (this-as self (emit self entry))))
  (unchecked-set proto "_onUpdate" (fn [entry] (this-as self (on-update self entry))))
  (unchecked-set proto "_sweepUnseen" (fn [] (this-as self (sweep-unseen self))))
  (unchecked-set proto "_authorOf" (fn [entry] (this-as self (author-of self entry))))
  (unchecked-set proto "_canJoinLocally" (fn [entry] (this-as self (can-join-locally self entry))))
  (unchecked-set proto "_identityLocal" (fn [e] (this-as self (identity-local self e)))))

static mirror of HEAD · about · clone: git clone https://git.ardegazu.ro/rooms-kit.git