vector-core 0.7.2

Core library for Vector — the single source of truth for all Vector clients, SDKs, and interfaces.
Documentation
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
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
//! v2 realtime — the `authors`-based live subscription + dispatch.
//!
//! v2 addresses planes by their group pubkey (CORD-01), not v1's `#z` tag, so
//! the subscription is `{kinds:[1059,21059], authors:[…plane pubkeys…]}`. A
//! received event is routed to the owning community and handed to the shared
//! [`inbound::dispatch_wrap`], which fires the protocol-agnostic
//! `InboundEventHandler` the SDK's `on_message` consumes.
//!
//! The kind-1059 dispatch rule (CLAUDE.md A4): the listen loop tries this v2
//! path for events on the v2 subscription; DM gift wraps and 3313 Direct Invites
//! (both `#p=me`) stay on the DM subscription — the author-set here never
//! includes an identity key, so the two never collide.

use std::collections::HashSet;
use std::sync::{Arc, LazyLock, Mutex as StdMutex};

use nostr_sdk::prelude::{Client, Event, Filter, Kind, PublicKey, RelayStatus, RelayUrl, SubscriptionId};
use tokio::sync::mpsc::UnboundedSender;
use tokio::sync::Mutex;

use super::community::CommunityV2;
use super::stream;
use super::{derive, inbound};
use crate::community::{CommunityId, ConcordProtocol, Epoch};
use crate::event_handler::InboundEventHandler;
use crate::state::SessionGuard;
use crate::ClientRelayExt;

/// The targeted subscription id (streams on desktop).
static V2_SUB_ID: LazyLock<Mutex<Option<SubscriptionId>>> = LazyLock::new(|| Mutex::new(None));
/// The pool-wide subscription id (the path that streams on Android).
static V2_POOLWIDE_SUB_ID: LazyLock<Mutex<Option<SubscriptionId>>> = LazyLock::new(|| Mutex::new(None));
/// The current author-set (sorted hex) — an unchanged set skips a churny
/// unsubscribe+resubscribe, exactly like v1's `COMMUNITY_SUB_SET`.
static V2_SUB_SET: LazyLock<Mutex<Vec<String>>> = LazyLock::new(|| Mutex::new(Vec::new()));

/// Session generation of the last COMPLETED [`refresh_subscription`] pass, or
/// `u64::MAX` before any. Generation-keyed rather than a bare bool so an account
/// swap invalidates it automatically — a per-account global that survived a swap
/// would report account A's readiness for account B.
static V2_SUB_READY_GEN: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(u64::MAX);

#[inline]
fn mark_subscription_ready() {
    V2_SUB_READY_GEN.store(
        crate::state::SessionGuard::capture().generation(),
        std::sync::atomic::Ordering::Release,
    );
}

/// Whether the v2 Community subscription pass has COMPLETED for the active
/// session: live channel events are being delivered (or there is genuinely
/// nothing to subscribe to). `false` during the cold-start window where the
/// relay-connect wait + stream-auth priming are still running — the state a
/// health-checking bot could previously not distinguish from "subscribed to
/// nothing".
pub fn subscription_ready() -> bool {
    V2_SUB_READY_GEN.load(std::sync::atomic::Ordering::Acquire)
        == crate::state::SessionGuard::capture().generation()
}
/// Outer-wrap ids already dispatched, so the handler fires EXACTLY ONCE per
/// message. The relay pool delivers the same wrap under both the targeted and
/// pool-wide subs and from every relay independently, so without this a bot's
/// `on_message` (and its reply) would run several times per message — the v1
/// community path and the DM path dedup by outer id for the same reason. Cleared
/// on session swap; coarsely bounded (a message's duplicates all arrive within a
/// short window, so a recent-set suffices).
static V2_SEEN_WRAPS: LazyLock<Mutex<HashSet<[u8; 32]>>> = LazyLock::new(|| Mutex::new(HashSet::new()));
/// Bound on [`V2_SEEN_WRAPS`] before a coarse flush.
const SEEN_WRAPS_CAP: usize = 8192;
/// The per-community follow QUEUE. dispatch, boot catch-up, reconnect, and manual
/// sync all just [`enqueue_follow`] — non-blocking + coalesced. A single spawned
/// worker ([`spawn_follow_worker`]) drains it and runs one combined rekey+control
/// follow per community at a time, so two triggers can never concurrently
/// whole-row-save and clobber each other. This replaces the old gate/rerun/
/// spawn-vs-await machinery: an enqueue never blocks its caller, and a junk-wrap
/// flood coalesces to at most one queued + one running follow per community.
/// Reset on session swap (the worker exits when its `SessionGuard` invalidates or
/// its channel closes).
static V2_FOLLOW_TX: LazyLock<StdMutex<Option<UnboundedSender<CommunityId>>>> = LazyLock::new(|| StdMutex::new(None));
/// Community ids currently queued or processing — coalesces a burst to one follow.
static V2_FOLLOW_PENDING: LazyLock<StdMutex<HashSet<[u8; 32]>>> = LazyLock::new(|| StdMutex::new(HashSet::new()));
/// Per-community follow serialization, shared by the queue worker AND the inline
/// (headless) follow path. The worker-vs-inline CHOICE is a benign race
/// (`follow_worker_running` is check-then-act; a worker can spawn right after a
/// `false`), so correctness can't ride on it: whichever path runs, the follow body
/// executes under this lock and two follows of one community can never interleave
/// their whole-row saves. Bounded by the held-community count; reset on swap.
static V2_FOLLOW_LOCKS: LazyLock<StdMutex<std::collections::HashMap<[u8; 32], Arc<Mutex<()>>>>> =
    LazyLock::new(|| StdMutex::new(std::collections::HashMap::new()));

/// The follow lock for one community (created on first use).
pub(crate) fn follow_lock(id: &CommunityId) -> Arc<Mutex<()>> {
    V2_FOLLOW_LOCKS.lock().unwrap().entry(id.0).or_default().clone()
}

/// The author set the live v2 subscription CURRENTLY carries (sorted hex).
/// Diagnostics: comparing this against the freshly-derived plane authors is the
/// one-glance test for "did a rotation leave this client subscribed to a dead
/// epoch" — the failure that otherwise only shows up as silent missing events.
pub async fn subscribed_author_set() -> Vec<String> {
    V2_SUB_SET.lock().await.clone()
}

/// Diagnostics: run the three follow stages once for `id` UNDER the follow lock
/// (so it can't race the worker into a whole-row clobber) and return each
/// stage's outcome as a display string. Reloads freshly-persisted state between
/// stages, exactly like [`follow_community`]. Mutating — same writes a live
/// follow makes; the lock serializes it against the worker.
#[cfg(debug_assertions)]
pub async fn debug_run_follow_stages(id: &CommunityId, session: &SessionGuard) -> (String, String, String) {
    let lock = follow_lock(id);
    let _guard = lock.lock().await;
    let transport = crate::community::transport::LiveTransport::with_timeout(std::time::Duration::from_secs(12));

    let Ok(Some(c)) = crate::db::community::load_community_v2(id) else {
        return ("community gone".into(), "-".into(), "-".into());
    };
    let rekeys = match super::service::follow_rekeys(&transport, &c, session).await {
        Ok(f) => format!(
            "Ok(updated={} self_removed={} dissolved={})",
            f.updated.as_ref().map(|u| format!("root_e{}", u.root_epoch.0)).unwrap_or_else(|| "no".into()),
            f.self_removed, f.dissolved
        ),
        Err(e) => format!("ERR: {e}"),
    };
    let Ok(Some(c)) = crate::db::community::load_community_v2(id) else {
        return (rekeys, "community gone".into(), "-".into());
    };
    let control = match super::service::follow_control(&transport, &c, session).await {
        Ok(v) => format!("Ok(changed={})", v.is_some()),
        Err(e) => format!("ERR: {e}"),
    };
    let Ok(Some(c)) = crate::db::community::load_community_v2(id) else {
        return (rekeys, control, "community gone".into());
    };
    let guestbook = match super::service::sync_guestbook(&transport, &c, session).await {
        Ok(fresh) => format!("Ok(fresh={})", fresh.len()),
        Err(e) => format!("ERR: {e}"),
    };
    (rekeys, control, guestbook)
}

pub async fn subscription_id() -> Option<SubscriptionId> {
    V2_SUB_ID.lock().await.clone()
}

pub async fn poolwide_subscription_id() -> Option<SubscriptionId> {
    V2_POOLWIDE_SUB_ID.lock().await.clone()
}

/// Clear the v2 realtime state on a session reset (called from `swap_session`
/// alongside v1's clear), so a stale sub id / author-set can't leak across accounts.
pub async fn clear() {
    *V2_SUB_ID.lock().await = None;
    *V2_POOLWIDE_SUB_ID.lock().await = None;
    V2_SUB_SET.lock().await.clear();
    V2_SEEN_WRAPS.lock().await.clear();
    // Drop the queue sender so the worker's channel closes and it exits (its
    // SessionGuard also invalidates); the next login spawns a fresh worker.
    *V2_FOLLOW_TX.lock().unwrap() = None;
    V2_FOLLOW_PENDING.lock().unwrap().clear();
    V2_FOLLOW_LOCKS.lock().unwrap().clear();
    // Account A's stream keys must not keep authenticating (or answering relay
    // challenges) once account B is live.
    super::streamauth::clear();
}

/// Every plane pubkey a set of v2 communities publishes under that
/// [`inbound::dispatch_wrap`] handles — the subscription author-set. Per
/// community: the guestbook, the control plane, and each channel's current
/// Chat-Plane address. Pure + deterministic (deduped, sorted) — the testable core.
///
/// The **control plane** rides here so a long-running bot follows metadata +
/// public-channel edits live ([`super::service::follow_control`] re-folds on a
/// recognized wrap), and the next-epoch rekey planes ride via [`rekey_authors`]
/// (each subscribed author has its `dispatch_wrap` arm — never one without the
/// other).
pub fn plane_authors(communities: &[CommunityV2]) -> Vec<PublicKey> {
    let mut out = Vec::new();
    for c in communities {
        out.push(derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch).pk());
        out.push(control_author(c));
        // The dissolved plane (CORD-02 §9) — so a mid-session dissolution seals live.
        out.push(super::derive::dissolved_group_key(c.id()).pk());
        for ch in &c.channels {
            // A KEYLESS private channel has no readable chat plane — channel_secret's
            // root fallback would subscribe the PUBLIC plane for it. Its rekey plane
            // (below) is still watched, which is how its key arrives.
            if ch.private && ch.key.is_none() {
                continue;
            }
            let (secret, epoch) = c.channel_secret(ch);
            out.push(derive::channel_group_key(&secret, &ch.id, epoch).pk());
        }
        out.extend(rekey_authors(c));
    }
    out.sort_by_key(|p| p.to_hex());
    out.dedup();
    out
}

/// This community's Control Plane address at the current root epoch — the single
/// source of truth shared by [`plane_authors`] (subscribe) and
/// [`inbound::dispatch_wrap`] (recognize) so the two can't drift.
pub(crate) fn control_author(c: &CommunityV2) -> PublicKey {
    derive::control_group_key(&c.community_root, c.id(), c.root_epoch).pk()
}

/// The next-epoch rekey plane addresses for a community: the base rotation
/// (`root_epoch + 1`) and each Private channel's rotation (`channel epoch + 1`),
/// both under the CURRENT community_root. This is the single source of truth for
/// which rekey wraps we subscribe AND recognize (`inbound::dispatch_wrap` calls
/// it), so the two can never drift. A Public channel has no independent rotation
/// (it rides the base), so only Private channels contribute a channel address.
pub(crate) fn rekey_authors(c: &CommunityV2) -> Vec<PublicKey> {
    // saturating: a bundle's epoch isn't covered by the community_id commitment, so
    // it's attacker-influenced — never let `epoch + 1` overflow (a u64::MAX epoch
    // just yields a dead address, never a panic).
    let mut out = vec![derive::base_rekey_group_key(&c.community_root, c.id(), Epoch(c.root_epoch.0.saturating_add(1))).pk()];
    for ch in &c.channels {
        if ch.private {
            out.push(derive::channel_rekey_group_key(&c.community_root, &ch.id, Epoch(ch.epoch.0.saturating_add(1))).pk());
        }
    }
    out
}

/// Load every locally-held, LIVE **v2** community (dispatching each id by its
/// stored protocol). The realtime layer folds these into the subscription +
/// routing, so excluding a DISSOLVED community here is what enforces CORD-02 §9
/// on the receive side: its chat/control/rekey planes stop being subscribed and
/// an arriving wrap for it is `NotOurs` (dropped, never honored). Held keys still
/// open old history through the explicit read paths — this only stops NEW events.
pub fn load_held_v2() -> Vec<CommunityV2> {
    let ids = crate::db::community::list_community_ids().unwrap_or_default();
    ids.iter()
        .filter(|id| matches!(crate::db::community::community_protocol(id).ok().flatten(), Some(ConcordProtocol::V2)))
        .filter(|id| !crate::db::community::get_community_dissolved(&crate::simd::hex::bytes_to_hex_32(&id.0)).unwrap_or(false))
        .filter_map(|id| crate::db::community::load_community_v2(id).ok().flatten())
        .collect()
}

/// Refresh the v2 subscription for the held communities: register
/// `{kinds:[1059,21059], authors:[…]}` on their relays (targeted + pool-wide,
/// mirroring v1). Idempotent on an unchanged author-set.
pub async fn refresh_subscription(client: &Client) {
    // Phase 1, LOCK-FREE: make sure the community relays are added + connected —
    // the slow part (a connect wait of up to ~6s). Holding the sub locks across
    // this stalled every concurrent dispatch/refresh behind one caller's connect.
    {
        let communities = load_held_v2();
        let mut relays: Vec<String> = communities.iter().flat_map(|c| c.relays.iter().cloned()).collect();
        relays.sort();
        relays.dedup();
        if !relays.is_empty() {
            // Community relays ride GOSSIP|PING (warm but excluded from pool-wide DM ops).
            for r in &relays {
                let _ = client.add_managed_relay(r.as_str()).capabilities(crate::community_relay_capabilities()).await;
            }
            client.connect().await;
            // Wait briefly for at least one relay to actually connect (a subscribe
            // against a still-connecting relay silently fails to register — same
            // trap as v1).
            let wanted: Vec<RelayUrl> = relays.iter().filter_map(|r| RelayUrl::parse(r).ok()).collect();
            for _ in 0..24 {
                let pool = client.relays().all().await;
                if wanted.iter().any(|u| pool.get(u).map(|r| r.status() == RelayStatus::Connected).unwrap_or(false)) {
                    break;
                }
                tokio::time::sleep(std::time::Duration::from_millis(250)).await;
            }
            // Register every held plane's key BEFORE subscribing, so the responder
            // can answer any NIP-42 challenge our REQs trigger. Cheap + local.
            for c in &communities {
                super::streamauth::register_community(c);
            }
        }
    }

    // Phase 2, LOCKED + bounded (no connect waits): RE-snapshot the held state
    // INSIDE the sub locks — the follow worker runs concurrently and may have just
    // adopted a rotation, so the LAST locker must read the freshest persisted
    // authors. An out-of-lock read lets a stale caller commit an old author-set
    // over a fresh one and silently mute a rotated community. (A community whose
    // relays appeared between the phases subscribes now and connects on the next
    // refresh — the follow that discovers it always triggers one.)
    let mut sub_guard = V2_SUB_ID.lock().await;
    let mut set_guard = V2_SUB_SET.lock().await;

    let communities = load_held_v2();
    let authors = plane_authors(&communities);
    let mut relays: Vec<String> = communities.iter().flat_map(|c| c.relays.iter().cloned()).collect();
    relays.sort();
    relays.dedup();

    let mut new_set: Vec<String> = authors.iter().map(|p| p.to_hex()).collect();
    new_set.sort();

    // Unchanged-set fast path — but only when BOTH subs actually registered (a
    // failed pool-wide subscribe would otherwise stay absent until the author-set
    // changes, and Android streams via the pool-wide path).
    if sub_guard.is_some() && *set_guard == new_set && (authors.is_empty() || V2_POOLWIDE_SUB_ID.lock().await.is_some()) {
        mark_subscription_ready(); // an unchanged live sub is a completed pass
        return; // the pool re-applies the live subs across reconnects.
    }
    if let Some(old) = sub_guard.take() {
        let _ = client.unsubscribe(&old).await;
    }
    *set_guard = new_set;

    if authors.is_empty() {
        if let Some(old_pw) = V2_POOLWIDE_SUB_ID.lock().await.take() {
            let _ = client.unsubscribe(&old_pw).await;
        }
        // Completed with nothing to hear: READY with zero planes, which is a
        // different observable state from "still connecting".
        mark_subscription_ready();
        return;
    }

    let filter = Filter::new()
        .kinds([Kind::Custom(stream::KIND_WRAP), Kind::Custom(stream::KIND_WRAP_EPHEMERAL)])
        .authors(authors)
        .limit(0);

    {
        let mut pw = V2_POOLWIDE_SUB_ID.lock().await;
        if let Some(old) = pw.take() {
            let _ = client.unsubscribe(&old).await;
        }
        if let Ok(out) = client.subscribe(filter.clone()).await {
            *pw = Some(out.value);
        }
    }
    if let Ok(out) = client
        .subscribe(nostr_sdk::prelude::ReqTarget::manual(
            relays.iter().cloned().map(|u| (u, vec![filter.clone()])),
        ))
        .await
    {
        *sub_guard = Some(out.value);
    }
    mark_subscription_ready();
    drop(set_guard);
    drop(sub_guard);

    // AUTH-gating relays serve the planes only to a stream-authenticated
    // connection, and a live subscription isn't auto-retried after the gate. The
    // priming used to run BEFORE the subscribe, which gated the whole registration
    // on the slowest relay: a set with dead or auth-refusing members left the
    // client deaf for a minute while healthy relays sat idle. Prime in the
    // BACKGROUND instead — healthy relays stream from the registration above, and
    // each gating relay's auth completion re-sends the subs (`resubscribe_relay`
    // via the responder), so it joins the moment IT is ready, gating nobody else.
    prime_auth_in_background(client, relays);
}

/// In-flight coalescer for the background auth prime: boot fires many refreshes
/// (dispatch, absorb, reconnect) and each spawning its own prime would stampede
/// the gating relays with duplicate auth fetches.
static PRIME_IN_FLIGHT: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);

fn prime_auth_in_background(client: &Client, relays: Vec<String>) {
    use std::sync::atomic::Ordering;
    if relays.is_empty() || PRIME_IN_FLIGHT.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire).is_err() {
        return;
    }
    let client = client.clone();
    let session = crate::state::SessionGuard::capture();
    tokio::spawn(async move {
        if session.is_valid() {
            super::streamauth::prime_auth(&client, &relays).await;
        }
        PRIME_IN_FLIGHT.store(false, Ordering::Release);
    });
}

/// Re-send the CURRENT v2 subscriptions (same ids) to ONE relay. An AUTH-gating
/// relay CLOSEs a sub REQ that raced ahead of the stream AUTHs on a fresh
/// connection, and it never re-challenges once the connection is authenticated —
/// so the moment the streams finish authenticating is exactly when the subs must
/// be re-sent. nostr-sdk re-sends them only when ITS OWN auth completes (it
/// can't see ours), which usually wins by socket ordering; this makes the heal
/// deterministic and independent of that internal. Same-id REQs are idempotent.
pub(crate) async fn resubscribe_relay(client: &Client, relay: &RelayUrl) {
    let targeted = V2_SUB_ID.lock().await.clone();
    let poolwide = V2_POOLWIDE_SUB_ID.lock().await.clone();
    if targeted.is_none() && poolwide.is_none() {
        return; // nothing subscribed yet — the first refresh registers on an authed socket.
    }
    let communities = load_held_v2();
    let authors = plane_authors(&communities);
    if authors.is_empty() {
        return;
    }
    let filter = Filter::new()
        .kinds([Kind::Custom(stream::KIND_WRAP), Kind::Custom(stream::KIND_WRAP_EPHEMERAL)])
        .authors(authors)
        .limit(0);
    for id in [targeted, poolwide].into_iter().flatten() {
        let _ = client
            .subscribe(nostr_sdk::prelude::ReqTarget::single(relay.clone(), [filter.clone()]))
            .with_id(id)
            .await;
    }
}

/// Route an arriving v2 wrap: find the held community whose plane it opens under
/// and fire the matching handler callback (via the shared bridge). Persistence to
/// the local DB is deferred (bots deliver via the callback; GUI history is v1 for
/// now). `session` gates against a mid-flight account swap.
pub async fn dispatch_event(session: &SessionGuard, event: Event, handler: Arc<dyn InboundEventHandler>) {
    let Some(my_pk) = crate::my_public_key() else {
        return;
    };
    if !session.is_valid() {
        return;
    }
    // Fire EXACTLY ONCE per wrap: the pool re-delivers the same event under both
    // subs and from every relay. `insert` returns false if already dispatched.
    {
        let mut seen = V2_SEEN_WRAPS.lock().await;
        if !seen.insert(event.id.to_bytes()) {
            return;
        }
        if seen.len() > SEEN_WRAPS_CAP {
            let keep = event.id.to_bytes();
            seen.clear();
            seen.insert(keep);
        }
    }
    let communities = load_held_v2();
    for c in &communities {
        match inbound::dispatch_wrap(&event, c, &my_pk, &*handler) {
            inbound::DispatchedV2::NotOurs => continue,
            // A control OR a rekey wrap: just enqueue a follow for this community.
            // Non-blocking + coalesced — the single follow worker serializes control
            // and rekey per community (no concurrent whole-row clobber) off this hot
            // path, so a junk-wrap flood can't head-of-line-block the notification loop.
            inbound::DispatchedV2::Control { .. } | inbound::DispatchedV2::Rekey { .. } => {
                enqueue_follow(c.id());
                return;
            }
            inbound::DispatchedV2::Dissolved { community_id } => {
                // Death wins (CORD-02 §9): seal read-only + surface the grave, ONCE
                // (a re-wrapped tombstone with a fresh outer id must not re-fire the
                // handler). The next load_held_v2 excludes it, so its planes also
                // stop being subscribed + routed.
                if crate::db::community::set_community_dissolved(&community_id).unwrap_or(false) {
                    handler.on_community_dissolved(&community_id);
                    if let Some(client) = crate::state::nostr_client() {
                        refresh_subscription(&client).await;
                    }
                }
                return;
            }
            // A chat event, opened but NOT yet applied: persist first (dedup by inner
            // id + the author-scoped edit/delete checks), then fire the callback from
            // the outcome — v1's exact model. A re-wrapped duplicate (any keyholder
            // can re-seal a signed rumor into a fresh 1059), the relay echo of our
            // own send, or a forged edit/delete yields no outcome and re-fires
            // nothing.
            inbound::DispatchedV2::Chat { channel_id, event } => {
                if !session.is_valid() {
                    return;
                }
                match inbound::persist_chat_event(&event, &channel_id, &my_pk, session).await {
                    Some(inbound::ChatPersist::New(message)) => handler.on_community_message(&channel_id, &message, true),
                    // A reaction or an edit: the folded TARGET row (its id is the
                    // target's) — the same payload v1 hands this callback. Both come
                    // back out of STATE, whose compact form drops the quoted
                    // attachment's extension, so re-resolve before rendering.
                    Some(inbound::ChatPersist::Updated { mut message, .. }) => {
                        let _ = crate::db::events::populate_reply_context(&mut message).await;
                        handler.on_community_update(&channel_id, &message.id, &message);
                    }
                    // An un-react re-renders the PARENT (its chips changed) — the
                    // same surface a landed reaction drives.
                    Some(inbound::ChatPersist::ReactionRemoved { mut message, .. }) => {
                        let _ = crate::db::events::populate_reply_context(&mut message).await;
                        handler.on_community_update(&channel_id, &message.id, &message);
                    }
                    Some(inbound::ChatPersist::Removed(target_id)) => handler.on_community_removed(&channel_id, &target_id),
                    None => {}
                }
                return;
            }
            inbound::DispatchedV2::Presence { .. } => {
                // Live membership motion: fold it into the persisted Guestbook so
                // the memberlist stays a local read (the presence callback already
                // fired inline). Reopen here — the dispatcher stays pure — and
                // refresh the overview when it lands.
                if !session.is_valid() {
                    return;
                }
                let gb = derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch);
                if let Ok(opened) = super::stream::open_wrap(&event, &gb) {
                    if let Ok(ev) = super::guestbook::parse_guestbook_event(&opened) {
                        let changed = super::service::ingest_guestbook_event(c, ev, event.created_at.as_secs()).unwrap_or(false);
                        if changed && session.is_valid() {
                            handler.on_community_refreshed(&crate::simd::hex::bytes_to_hex_32(&c.id().0));
                        }
                    }
                }
                return;
            }
            inbound::DispatchedV2::Kick { target } => {
                if !session.is_valid() {
                    return;
                }
                let gb = derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch);
                let Ok(opened) = super::stream::open_wrap(&event, &gb) else { return };
                let Ok(ev) = super::guestbook::parse_guestbook_event(&opened) else { return };
                if !super::service::ingest_guestbook_event(c, ev, event.created_at.as_secs()).unwrap_or(false) {
                    return;
                }
                if !session.is_valid() {
                    return;
                }
                let community_id = crate::simd::hex::bytes_to_hex_32(&c.id().0);
                // Self-removal rides the AUTHORIZED fold, never the wrap: the store takes
                // any Kick rumor, so acting on arrival would let any member evict anyone.
                // Fail-closed — a kick our roster can't justify leaves us in place.
                //
                // The COALESCE verdict, not the memberlist: on a rejoin the store starts
                // empty while the control fold has already re-derived our old ban mark, so
                // the memberlist legitimately excludes us for the window before the
                // Guestbook catches up. A stale Kick landing in that window would read as
                // an authorized eviction. Coalescing asks only what our latest authorized
                // entry is, so a fresh Join beats the old Kick and an empty store decides
                // nothing.
                let evicted = my_pk == target && super::service::stored_kick_verdict(c, &my_pk);
                if evicted {
                    // WARN, not info: `log_info!` is compiled out in release, and a
                    // self-eviction is exactly the destructive, hard-to-diagnose event
                    // that must leave a trace on a shipped build.
                    crate::log_warn!(
                        "[v2:teardown {}] KICK: the authorized guestbook fold rules us kicked",
                        &community_id[..8.min(community_id.len())]
                    );
                    handler.on_community_self_removed(&community_id);
                } else {
                    crate::log_debug!(
                        "[v2:kick {}] declined: target={} is not ruled kicked by the fold",
                        &community_id[..8.min(community_id.len())], &target.to_hex()[..8]
                    );
                    handler.on_community_refreshed(&community_id);
                }
                return;
            }
            _ => return, // typing (and non-surfaced guestbook kinds) handled inline by the dispatcher.
        }
    }
}

/// Whether a live follow worker is draining the queue (a `listen()` is running).
/// Headless callers use this to run a follow inline instead of enqueueing into
/// the void.
pub fn follow_worker_running() -> bool {
    V2_FOLLOW_TX.lock().unwrap().as_ref().map(|tx| !tx.is_closed()).unwrap_or(false)
}

/// Queue a follow for `id` — NON-BLOCKING + coalesced. A burst (or a junk-wrap
/// flood) collapses to at most one queued + one running follow per community. A
/// no-op if no worker is running (no live `listen()`). Callers: dispatch,
/// boot/reconnect catch-up, manual sync — none of them block or touch a lock for
/// longer than the enqueue.
pub fn enqueue_follow(id: &CommunityId) {
    let mut pending = V2_FOLLOW_PENDING.lock().unwrap();
    if !pending.insert(id.0) {
        return; // already queued or processing — coalesce.
    }
    match V2_FOLLOW_TX.lock().unwrap().as_ref() {
        Some(tx) if tx.send(*id).is_ok() => {}
        _ => {
            // No worker yet — PARK the id in the pending set instead of
            // dropping it; spawn_follow_worker drains parked ids into its
            // fresh queue. The boot sweep's enqueues race notifs' worker
            // spawn (they fire the moment init completes), so a pre-worker
            // enqueue must defer, never silently vanish — a dropped boot
            // refold leaves an offline-rotated community wedged at its old
            // epoch until some live event happens to trigger a dispatch.
        }
    }
}

/// Spawn the single follow worker for this session. Installs the queue sender and
/// drains it, running one combined follow per community at a time. Replacing the
/// sender (a re-`listen()`) or [`clear`] (a swap) closes the old channel so the old
/// worker exits; the captured `SessionGuard` also stops it. Idempotent per session.
pub fn spawn_follow_worker(handler: Arc<dyn InboundEventHandler>) {
    let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<CommunityId>();
    // Re-send every pending id into the fresh queue: parked pre-worker
    // enqueues AND anything a replaced worker's dropped channel still owed.
    // Entries stay in the set — the worker removes each as it starts
    // processing, preserving the coalescing invariant.
    {
        let pending = V2_FOLLOW_PENDING.lock().unwrap();
        for id in pending.iter() {
            let _ = tx.send(CommunityId(*id));
        }
    }
    *V2_FOLLOW_TX.lock().unwrap() = Some(tx);
    let session = SessionGuard::capture();
    tokio::spawn(async move {
        while let Some(id) = rx.recv().await {
            if !session.is_valid() {
                break;
            }
            // Remove from pending BEFORE running, so a trigger arriving DURING the
            // follow re-enqueues (and is processed after) rather than being lost.
            V2_FOLLOW_PENDING.lock().unwrap().remove(&id.0);
            follow_community(&session, &id, &*handler).await;
        }
    });
}

/// One combined rekey-then-control follow for a community, each pass against the
/// FRESHLY-RELOADED persisted state (never a stale clone, so the two planes can't
/// lose each other's writes). Rekey runs first: a base adopt moves the control
/// address, and a self-removal tears the community down (skipping control). No-op
/// without a live client — unit tests drive `service::follow_control` /
/// `follow_rekeys` directly.
async fn follow_community(session: &SessionGuard, id: &CommunityId, handler: &dyn InboundEventHandler) {
    let Some(client) = crate::state::nostr_client() else {
        return;
    };
    // Serialize against an inline (headless) follow of the same community — the
    // queue only serializes triggers routed THROUGH it.
    let lock = follow_lock(id);
    let _guard = lock.lock().await;
    let community_id = crate::simd::hex::bytes_to_hex_32(&id.0);
    let transport = crate::community::transport::LiveTransport::with_timeout(std::time::Duration::from_secs(12));

    // Rekey first (fresh DB state).
    let Ok(Some(current)) = crate::db::community::load_community_v2(id) else {
        return; // community gone (left / removed).
    };
    match super::service::follow_rekeys(&transport, &current, session).await {
        // A tombstone surfaced during catch-up (an offline member learning of a
        // death) — the flag is set; seal + surface, and stop following.
        Ok(follow) if follow.dissolved => {
            if !session.is_valid() {
                return;
            }
            crate::log_warn!("[v2:teardown {}] DISSOLVED tombstone", &community_id[..8.min(community_id.len())]);
            handler.on_community_dissolved(&community_id);
            return;
        }
        Ok(follow) if follow.self_removed => {
            if !session.is_valid() {
                return;
            }
            crate::log_warn!("[v2:teardown {}] REKEY EXCLUSION: an authorized rotation left us out", &community_id[..8.min(community_id.len())]);
            let _ = crate::db::community::delete_community(&community_id);
            refresh_subscription(&client).await;
            handler.on_community_self_removed(&community_id);
            return;
        }
        Ok(follow) if follow.updated.is_some() => {
            if !session.is_valid() {
                return;
            }
            refresh_subscription(&client).await;
            handler.on_community_refreshed(&community_id);
        }
        Ok(_) => {}
        Err(e) => {
            // Surfaced, not swallowed: with EOSE-verified fetches an AUTH-gated
            // or dead rekey plane now reports here instead of masquerading as
            // "no rotation" — the exact signature of an epoch wedge.
            crate::log_warn!("[v2:follow {}] rekey follow failed (will retry on next trigger): {}", &community_id[..8.min(community_id.len())], e);
            return;
        }
    }

    // Control second, on the (possibly new-root) freshly-reloaded state.
    let Ok(Some(current)) = crate::db::community::load_community_v2(id) else {
        return;
    };
    let control_changed = matches!(
        super::service::follow_control(&transport, &current, session).await,
        Ok(Some(_))
    );
    if !session.is_valid() {
        return;
    }
    // Re-judge parked key vends on EVERY pass, not only when the fold moved. A
    // vend that lands AFTER its Grant folded has nothing left to change the
    // control plane, so gating this on `changed` strands it until some unrelated
    // edition happens by (CORD-03 "delivered on grant").
    if let Ok(Some(folded)) = crate::db::community::load_community_v2(id) {
        let keyed = super::service::absorb_parked_channel_keys(&folded, session);
        if !keyed.is_empty() {
            crate::log_info!(
                "[v2:follow {}] adopted {} vended private-channel key(s)",
                &community_id[..8.min(community_id.len())],
                keyed.len()
            );
            // A newly-keyed channel adds its chat plane to the readable set, and
            // the live subscription is addressed BY that set — without a rebuild
            // we hold the key and still never hear the channel. Independent of
            // `control_changed`: adopting a parked vend changes no control state.
            refresh_subscription(&client).await;
        }
        for ch in keyed {
            let hex = crate::simd::hex::bytes_to_hex_32(&ch.0);
            // Read what was said while we were locked out. The live subscription
            // only carries what arrives NEXT, so without this a member granted
            // access walks into a room that looks empty. The count rides the hook:
            // backfilled history is NOT dispatched as live messages, and without
            // the number a handler cannot even tell there is history to go read.
            let backfilled = crate::VectorCore::v2_backfill_channel(
                id, &hex, 50, 2, None, None,
                crate::community::transport::Evidence::Fast, 12,
            )
            .await;
            if !session.is_valid() {
                return;
            }
            handler.on_channel_keyed(&community_id, &hex, backfilled);
        }
    }
    if control_changed {
        refresh_subscription(&client).await;
        handler.on_community_refreshed(&community_id);
        // A control change can reveal rekey work that predates it — a just-announced
        // private channel's key crate is already sitting on its rekey plane (the key
        // ships BEFORE the vsk-2), and this pass's rekey walk ran before the channel
        // existed. Queue one more pass; it coalesces and converges (an unchanged
        // control fold doesn't re-queue).
        enqueue_follow(id);
    }

    // Banned by the banlist we just folded → self-remove NOW, not when the Refounding
    // eventually lands. A Ban composes three layers (CORD-04 §6) and their guarantees
    // arrive at different speeds: the Banlist edition is instant, the Refounding is
    // "heavy and asynchronous" by design. Waiting on the rotation to notice our own
    // removal left the target sitting in a community that had already dropped them from
    // its member count — silenced, still looking joined. Checked unconditionally: a
    // banlist-only change folds no document, so `follow_control` returns None for it.
    // v1 has done this since it shipped (`am_i_banned` in its own realtime refresh).
    if let Some(me) = crate::my_public_key() {
        if crate::db::community::is_author_banned(&community_id, &me) {
            if !session.is_valid() {
                return;
            }
            crate::log_warn!("[v2:teardown {}] SELF-BAN: our npub is in the folded banlist", &community_id[..8.min(community_id.len())]);
            handler.on_community_self_removed(&community_id);
            return;
        }
    }

    // Guestbook third: catch the membership store up from its cursor. Boot and
    // reconnect land here through this same queue, so the memberlist is a local
    // read by the time any panel asks — and every join/leave the catch-up folds
    // surfaces as a presence line (real-rumor-id keyed, so a line each path also
    // saw live inserts exactly once).
    let Ok(Some(current)) = crate::db::community::load_community_v2(id) else {
        return;
    };
    if let Ok(fresh) = super::service::sync_guestbook(&transport, &current, session).await {
        if fresh.is_empty() || !session.is_valid() {
            return;
        }
        surface_presence(&current, &fresh, handler);
        handler.on_community_refreshed(&community_id);
    }
}

/// Fire the presence-line surface for freshly-folded guestbook events — the
/// catch-up twin of the live dispatch (same handler, same real-id dedup key,
/// same banned-author drop). Kicks/snapshots shape the memberlist, not the feed.
fn surface_presence(
    community: &CommunityV2,
    fresh: &[super::guestbook::GuestbookEvent],
    handler: &dyn InboundEventHandler,
) {
    use super::guestbook::GuestbookEntry;
    use nostr_sdk::prelude::ToBech32;
    let Some(primary) = community.primary_channel() else {
        return;
    };
    let chat_id = crate::simd::hex::bytes_to_hex_32(&primary.id.0);
    let cid_hex = crate::simd::hex::bytes_to_hex_32(&community.id().0);
    let banned = crate::db::community::banned_set(&cid_hex);
    for ev in fresh {
        let (member, joined, at_ms, invited_by) = match &ev.entry {
            GuestbookEntry::Join { member, at_ms, invited_by } => (member, true, *at_ms, invited_by.clone()),
            GuestbookEntry::Leave { member, at_ms } => (member, false, *at_ms, None),
            GuestbookEntry::Kick { .. } | GuestbookEntry::Snapshot { .. } => continue,
        };
        if banned.contains(&member.to_bytes()) {
            continue;
        }
        let Ok(npub) = member.to_bech32();
        let event_id = crate::simd::hex::bytes_to_hex_32(&ev.rumor_id);
        let (by, label) = match &invited_by {
            Some((c, l)) => (Some(c.as_str()), Some(l.as_str())),
            None => (None, None),
        };
        handler.on_community_presence(&chat_id, &npub, joined, &event_id, at_ms / 1000, by, label);
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use super::super::control::{genesis, CommunityMetadata};
    use crate::community::Epoch;
    use nostr_sdk::prelude::Keys;

    fn a_community(name: &str) -> CommunityV2 {
        let owner = Keys::generate();
        let g = genesis(&owner, CommunityMetadata { name: name.into(), ..Default::default() }, 1_000).unwrap();
        CommunityV2::from_genesis(&g, name, None, vec!["wss://r".into()], 0)
    }

    #[test]
    fn subscription_readiness_is_scoped_to_the_session_that_marked_it() {
        // Generation-keyed on purpose: a bare bool would survive an account swap
        // and report account A's readiness for account B — the exact class of
        // per-account-global bug the swap-cache sweep exists to prevent.
        //
        // Bumping the GLOBAL generation invalidates every live SessionGuard, so
        // serialize with the swap-sensitive tests via the same guard they hold.
        let _guard = crate::db::DB_TEST_GUARD.lock().unwrap_or_else(|e| e.into_inner());
        mark_subscription_ready();
        assert!(subscription_ready(), "marked for the current session");
        crate::state::bump_session_generation();
        assert!(!subscription_ready(), "an account swap invalidates it");
        mark_subscription_ready();
        assert!(subscription_ready(), "the new session's own pass re-arms it");
    }

    #[test]
    fn plane_authors_covers_the_dispatched_planes_only() {
        let c = a_community("A");
        let authors = plane_authors(std::slice::from_ref(&c));

        // Subscribed: the guestbook, the control plane, and the one public channel
        // (the planes dispatch_wrap handles). Exactly those three.
        let gb = derive::guestbook_group_key(&c.community_root, c.id(), c.root_epoch).pk();
        let control = derive::control_group_key(&c.community_root, c.id(), c.root_epoch).pk();
        let general = {
            let (s, e) = c.channel_secret(&c.channels[0]);
            derive::channel_group_key(&s, &c.channels[0].id, e).pk()
        };
        // Plus the next base-rekey address (rekey-follow). A public channel has no
        // independent rotation, so #general contributes no channel-rekey address.
        let next_base = derive::base_rekey_group_key(&c.community_root, c.id(), Epoch(1)).pk();
        // Plus the dissolved plane (CORD-02 §9), so a live dissolution is detected.
        let dissolved = derive::dissolved_group_key(c.id()).pk();
        assert!(
            authors.contains(&gb)
                && authors.contains(&control)
                && authors.contains(&general)
                && authors.contains(&next_base)
                && authors.contains(&dissolved)
        );
        assert_eq!(authors.len(), 5, "guestbook + control + dissolved + chat + base-rekey planes are subscribed");
    }

    #[test]
    fn plane_authors_is_deterministic_deduped_and_multi_community() {
        let a = a_community("A");
        let b = a_community("B");
        let one = plane_authors(std::slice::from_ref(&a));
        // Re-running over the same community is byte-identical (deterministic).
        assert_eq!(plane_authors(std::slice::from_ref(&a)), one);
        // Two distinct communities' planes are all present, none dropped.
        let two = plane_authors(&[a.clone(), b.clone()]);
        assert_eq!(two.len(), one.len() * 2);
        // Order-independent: reversing the input yields the identical sorted set.
        assert_eq!(plane_authors(&[b, a]), two);
    }

    #[tokio::test]
    async fn dispatch_event_routes_a_v2_message_to_the_handler() {
        use crate::community::transport::memory::MemoryRelay;
        use crate::community::transport::{Query, Transport};
        use crate::types::Message;
        use std::sync::Mutex as StdMutex;

        #[derive(Default)]
        struct Recorder {
            got: StdMutex<Vec<(String, String)>>,
        }
        impl InboundEventHandler for Recorder {
            fn on_community_message(&self, chat_id: &str, msg: &Message, _new: bool) {
                self.got.lock().unwrap().push((chat_id.to_string(), msg.content.clone()));
            }
        }

        // Offline DB + identity.
        let _g = crate::db::DB_TEST_GUARD.lock().unwrap_or_else(|e| e.into_inner());
        crate::db::close_database();
        crate::db::clear_id_caches();
        let tmp = tempfile::tempdir().unwrap();
        let acct = {
            const B: &[u8] = b"qpzry9x8gf2tvdw0s3jn54khce6mua7l";
            let mut s = String::from("npub1");
            for i in 0..58 {
                s.push(B[(i * 5 + 1) % 32] as char);
            }
            s
        };
        std::fs::create_dir_all(tmp.path().join(&acct)).unwrap();
        crate::db::set_app_data_dir(tmp.path().to_path_buf());
        crate::db::set_current_account(acct.clone()).unwrap();
        crate::db::init_database(&acct).unwrap();
        let _ = crate::state::take_nostr_client();
        let me = Keys::generate();
        crate::state::MY_SECRET_KEY.store_from_keys(&me, &[]);
        crate::state::set_my_public_key(me.public_key());

        // Create a v2 community (persisted), then ANOTHER member (holds the root)
        // posts — the incoming case a live sub delivers (an OWN send is echoed at
        // send time, so its relay copy correctly dedups instead of firing).
        let relay = MemoryRelay::new();
        let community = super::super::service::create_community(&relay, "Live", vec!["wss://r".into()], None).await.unwrap();
        let general = community.channels[0].id;
        let member = Keys::generate();
        let group = derive::channel_group_key(&community.community_root, &general, community.root_epoch);
        let rumor = super::super::chat::build_message_rumor(member.public_key(), &general, community.root_epoch, "live ping", None, &[], vec![], 5_000);
        let (wrap, _) = super::super::chat::seal_chat_rumor(&rumor, &group, &member, nostr_sdk::prelude::Timestamp::from_secs(5), false).unwrap();
        let _ = relay.publish(&wrap, &community.relays).await;
        let q = Query { kinds: vec![stream::KIND_WRAP], authors: vec![group.pk_hex()], ..Default::default() };
        let wrap = relay.fetch(&q, &community.relays).await.unwrap().into_iter().find(|w| w.pubkey == group.pk()).unwrap();

        // The realtime dispatch (loading held v2 communities from the DB) routes it.
        // Dispatch the SAME wrap TWICE — modelling the pool re-delivering it under
        // the targeted + pool-wide subs (and from multiple relays). The handler
        // must fire EXACTLY ONCE (no duplicate bot replies) — and only AFTER the
        // persist outcome (the callbacks-from-persist model).
        let rec = Arc::new(Recorder::default());
        let session = SessionGuard::capture();
        crate::community::v2::realtime::clear().await; // fresh seen-set for the test
        dispatch_event(&session, wrap.clone(), rec.clone()).await;
        dispatch_event(&session, wrap, rec.clone()).await;

        let got = rec.got.lock().unwrap();
        assert_eq!(got.len(), 1, "a re-delivered wrap fires the handler exactly once");
        assert_eq!(got[0].1, "live ping");
        assert_eq!(got[0].0, crate::simd::hex::bytes_to_hex_32(&general.0));
    }

    #[tokio::test]
    async fn follow_queue_coalesces_a_burst_and_re_enqueues_after_processing() {
        // Install a test channel in place of the worker's, so we can observe what the
        // queue delivers without a live client.
        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<CommunityId>();
        *V2_FOLLOW_TX.lock().unwrap() = Some(tx);
        V2_FOLLOW_PENDING.lock().unwrap().clear();

        let id = CommunityId([0x11; 32]);
        // A burst for one community collapses to a SINGLE queued follow (coalesced).
        enqueue_follow(&id);
        enqueue_follow(&id);
        enqueue_follow(&id);
        assert_eq!(rx.recv().await, Some(id), "first trigger queues a follow");
        assert!(rx.try_recv().is_err(), "the burst coalesced to exactly one");

        // The worker removes it from pending before running; a trigger AFTER that
        // re-queues (so a change during a follow isn't lost).
        V2_FOLLOW_PENDING.lock().unwrap().remove(&id.0);
        enqueue_follow(&id);
        assert_eq!(rx.recv().await, Some(id), "a trigger after processing re-queues");

        // A different community is independent (not coalesced against the first).
        let id2 = CommunityId([0x22; 32]);
        enqueue_follow(&id2);
        assert_eq!(rx.recv().await, Some(id2));

        *V2_FOLLOW_TX.lock().unwrap() = None;
        V2_FOLLOW_PENDING.lock().unwrap().clear();
    }

    #[tokio::test]
    async fn a_dissolved_community_honors_no_new_events_and_fires_death_once() {
        use super::super::service;
        use crate::community::transport::memory::MemoryRelay;
        use crate::community::transport::Transport;
        use crate::types::Message;
        use std::sync::Mutex as StdMutex;

        #[derive(Default)]
        struct Recorder {
            messages: StdMutex<Vec<String>>,
            deaths: StdMutex<Vec<String>>,
        }
        impl InboundEventHandler for Recorder {
            fn on_community_message(&self, _chat: &str, msg: &Message, _new: bool) {
                self.messages.lock().unwrap().push(msg.content.clone());
            }
            fn on_community_dissolved(&self, community_id: &str) {
                self.deaths.lock().unwrap().push(community_id.to_string());
            }
        }

        let _g = crate::db::DB_TEST_GUARD.lock().unwrap_or_else(|e| e.into_inner());
        crate::db::close_database();
        crate::db::clear_id_caches();
        let tmp = tempfile::tempdir().unwrap();
        let acct = {
            const B: &[u8] = b"qpzry9x8gf2tvdw0s3jn54khce6mua7l";
            let mut s = String::from("npub1");
            for i in 0..58 {
                s.push(B[(i * 3 + 2) % 32] as char);
            }
            s
        };
        std::fs::create_dir_all(tmp.path().join(&acct)).unwrap();
        crate::db::set_app_data_dir(tmp.path().to_path_buf());
        crate::db::set_current_account(acct.clone()).unwrap();
        crate::db::init_database(&acct).unwrap();
        let _ = crate::state::take_nostr_client();
        let me = Keys::generate();
        crate::state::MY_SECRET_KEY.store_from_keys(&me, &[]);
        crate::state::set_my_public_key(me.public_key());

        let relay = MemoryRelay::new();
        let community = service::create_community(&relay, "Doomed", vec!["wss://r".into()], None).await.unwrap();
        let general = community.channels[0].id;

        // The owner's tombstone arrives as a MEMBER sees it (local flag still 0 —
        // build + publish it directly rather than via dissolve_community, which
        // would seal our own DB first and make it the owner-published case). `me`
        // is the owner, so the seal verifies.
        let rumor = super::super::dissolution::dissolved_tombstone_rumor(me.public_key(), community.id(), 8_000);
        let tombstone = super::super::dissolution::seal_dissolved(&rumor, community.id(), &me, nostr_sdk::prelude::Timestamp::from_secs(8_000)).unwrap();
        let _ = relay.publish(&tombstone, &community.relays).await;
        assert!(!crate::db::community::get_community_dissolved(&crate::simd::hex::bytes_to_hex_32(&community.id().0)).unwrap(), "not yet locally sealed");

        let rec = Arc::new(Recorder::default());
        let session = SessionGuard::capture();
        clear().await;
        // Fire the SAME tombstone twice AND a fresh re-wrap of its verified seal
        // (distinct outer id) — death must be announced exactly once.
        dispatch_event(&session, tombstone.clone(), rec.clone()).await;
        dispatch_event(&session, tombstone, rec.clone()).await;
        assert!(crate::db::community::get_community_dissolved(&crate::simd::hex::bytes_to_hex_32(&community.id().0)).unwrap());

        // A member posts a fresh message to #general AFTER the tombstone. It must
        // not be honored (the community is excluded from load_held_v2, so its
        // plane is NotOurs).
        let member = Keys::generate();
        let cgroup = derive::channel_group_key(&community.community_root, &general, community.root_epoch);
        let rumor = super::super::chat::build_message_rumor(member.public_key(), &general, community.root_epoch, "into the grave", None, &[], vec![], 9_000);
        let (mw, _) = super::super::chat::seal_chat_rumor(&rumor, &cgroup, &member, nostr_sdk::prelude::Timestamp::from_secs(9), false).unwrap();
        dispatch_event(&session, mw, rec.clone()).await;

        assert_eq!(rec.deaths.lock().unwrap().len(), 1, "death is announced exactly once");
        assert!(rec.messages.lock().unwrap().is_empty(), "a post-tombstone message is never honored (CORD-02 §9)");
    }

    #[test]
    fn a_private_channel_subscribes_to_its_own_chat_plane() {
        let mut c = a_community("Priv");
        c.channels.push(super::super::community::ChannelV2 {
            id: crate::community::ChannelId([0x33; 32]),
            name: "mods".into(),
            private: true,
            key: Some([0x44; 32]),
            epoch: Epoch(1),
            voice: None,
            meta_custom: None,
            meta_extra: Default::default(),
        });
        let authors = plane_authors(std::slice::from_ref(&c));
        // A private channel is read under its OWN key/epoch (not the root).
        let priv_chat = derive::channel_group_key(&[0x44; 32], &c.channels[1].id, Epoch(1)).pk();
        assert!(authors.contains(&priv_chat), "a private channel subscribes to its own chat plane");
        // Its next-rekey address IS subscribed (rekey-follow), keyed by the current
        // root at the channel's next epoch — so a rotation is delivered.
        let next_rekey = derive::channel_rekey_group_key(&c.community_root, &c.channels[1].id, Epoch(2)).pk();
        assert!(authors.contains(&next_rekey), "a private channel's next rekey plane is subscribed");
    }
}