From de5fa6625aa5693afcde7335a572456680aef9df Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Fri, 14 Aug 2026 17:09:42 +0000 Subject: [PATCH] feat(sync): probe participant NIP-65 mailboxes Repository conversations can continue in the read or write relays advertised by authors of accepted replies, reactions, zaps, and other descendants. Root-author inbox discovery alone therefore leaves valid indirect descendants undiscovered. Carry exact accepted-root provenance through the existing bounded descendant frontier, admit identity events only for those participants, and derive read/write mailbox ownership from accepted kind 10002 events. Probe one byte-bounded historic filter at a time with the existing per-relay fetch_events pagination, pacing, ledger, timeout, policy, and persistence paths. Starts are paced globally, but progress and terminal state remain independent per relay. Require both an established socket and the sync actor committed lifecycle before starting a probe. Prefer lifecycle-active due relays while falling back to the existing oldest-due dial order, so hundreds of unavailable mailbox sources cannot starve already-ready work. Drain completed mailbox probes before accepting more connection results so a busy startup queue cannot delay cursor progress or resource release. Production canaries exposed these startup conditions without requiring cross-relay coordination. The recursive descendant limit bounds which indirect IDs remain query roots; direct root references remain complete. Mailbox filter cursors are intentionally best-effort in-memory state: restart reconstructs ownership from LMDB and safely begins historic coverage again. This does not add permanent participant live subscriptions or a cross-relay completion coordinator. Validated with cargo check, strict all-target Clippy, 767 library tests including lifecycle and ready-selection regressions, and the three proactive Sync+ integration scenarios, including a write-mailbox child that references only a participant reaction. --- docs/explanation/README.md | 4 +- docs/explanation/architecture.md | 11 +- .../grasp-03-proactive-sync-plus.md | 118 ++- docs/explanation/peer-controlled-state.md | 2 +- src/nostr/builder.rs | 20 +- src/nostr/policy/identity.rs | 16 +- src/nostr/policy/mod.rs | 4 +- src/sync/discovery.rs | 85 ++ src/sync/mod.rs | 915 ++++++++++++++++-- src/sync/self_subscriber.rs | 8 +- tests/sync/proactive_sync_plus.rs | 114 ++- 11 files changed, 1157 insertions(+), 140 deletions(-) diff --git a/docs/explanation/README.md b/docs/explanation/README.md index d7bd6ca..6233798 100644 --- a/docs/explanation/README.md +++ b/docs/explanation/README.md @@ -109,9 +109,9 @@ Explanation documentation helps you **understand concepts** and design decisions --- ### [GRASP-03 Proactive Sync Plus](grasp-03-proactive-sync-plus.md) -**NIP-65 inbox discovery for accepted repository conversations** +**NIP-65 inbox/outbox discovery for accepted repository conversations** -**Read when:** You want to understand how accepted conversations are recovered from root-author inboxes +**Read when:** You want to understand how accepted conversations are recovered from participant mailboxes --- diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 12ad859..bac5739 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -593,8 +593,15 @@ For full design details, see [grasp-02-proactive-sync.md](grasp-02-proactive-syn GRASP-03 Sync+ is a default-on mailbox-discovery overlay on this manager. Set `NGIT_SYNC_PLUS_ENABLED=false` to retain GRASP-02 sync without discovering -root-author NIP-65 inboxes; NIP-11 advertises `GRASP-03` only while the overlay -is enabled. +accepted root and descendant participants' NIP-65 read/write mailboxes; NIP-11 +advertises `GRASP-03` only while the overlay is enabled. Exact thread +provenance scopes each mailbox to its accepted roots. Accepted root authors' +read/unmarked inboxes retain ordinary live/rotating GRASP-02 coverage. Wider +participant mailboxes use independent, history-only `fetch_events` workers: +at most one per relay and one new start per maintenance pass. They reuse the +connection's ordinary pacing, subscription ledger, pagination and 30-second +per-page terminal timeout without installing permanent participant +subscriptions or coupling progress between relays. ### Rejected Events Index diff --git a/docs/explanation/grasp-03-proactive-sync-plus.md b/docs/explanation/grasp-03-proactive-sync-plus.md index a53b746..58caaae 100644 --- a/docs/explanation/grasp-03-proactive-sync-plus.md +++ b/docs/explanation/grasp-03-proactive-sync-plus.md @@ -1,8 +1,9 @@ # GRASP-03 proactive sync plus GRASP-03 extends repository-declared GRASP-02 coverage to the Nostr outbox -model. An accepted issue, patch or pull request may have replies in the root -author's inbox even when those events are absent from repository relays. +model. An accepted issue, patch or pull request may have replies, reactions or +zaps in a conversation participant's inbox or outbox even when those events +are absent from repository relays. The overlay is enabled by default and can be disabled with `NGIT_SYNC_PLUS_ENABLED=false`. Because it extends the proactive GRASP-02 @@ -12,14 +13,19 @@ manager, it is effective only while that manager is running. NIP-11 advertises ## Minimal approach The implementation adds a narrow discovery control-plane, not a second sync -engine: derive accepted root IDs and authors locally; ask a small existing -index set for each author's latest profile kind `0` and NIP-65 kind `10002`; -retain those replaceable events locally; treat `read` and unmarked relays as -inboxes; and merge inbox-to-root mappings into ordinary per-relay sync targets. +engine: derive accepted root IDs and participant authors locally; ask a small +existing index set for each author's latest profile kind `0` and NIP-65 kind +`10002`; retain those replaceable events locally; keep accepted root authors' +`read` and unmarked inboxes in ordinary GRASP-02 coverage; and probe accepted +participants' `read`, `write` and unmarked conversation mailboxes with +repository and thread-reference filters. -GRASP-02 then owns all transport. Its current root tiers, historic and -live/rotating coverage, grouping, subscription ledger, rate-limit recovery, -pagination and per-relay connection lifecycle apply unchanged. +Mailbox work reuses GRASP-02's connection safety, filter byte packing, +pagination, subscription ledger, request pacing, event pipeline and write +policy. The participant expansion is history-only: a non-root participant +relay never becomes a permanent live source merely because an accepted author +advertised it. Root-author inboxes retain the pre-existing live/rotating Sync+ +behavior. Authors without an accepted stored relay list are queried through the configured user-index set plus the bootstrap relay. Once a list is retained, discovery @@ -28,21 +34,26 @@ to converge without making arbitrary repository relays identity sources. Discovery remains best effort and bounded to this operator-visible source graph. If a successful user-index query returns no accepted kind `10002`, accepted -roots by that author temporarily use the operator-configured Sync+ fallback -relay set as inboxes. This is the same root-only overlay consumed by ordinary -GRASP-02 historic, live and rotating coverage; it creates no separate -subscription scheduler. An accepted relay list removes the author from desired -fallback coverage and installs its declared inboxes. Existing shared live -subscriptions still drain naturally, as with any other relay-list replacement. +roots associated with that author temporarily use the operator-configured +Sync+ fallback relay set as mailboxes. Root-author fallback roots continue +through ordinary GRASP-02 historic, live and rotating coverage; non-root +participant roots use the paced history probe. An accepted relay list removes +the author from fallback coverage and schedules its declared relays instead. The self-subscriber builds a compact root-candidate inventory during its existing startup load and maintains it incrementally as roots arrive. StateOnly candidates remain inert; promotion of their repository to Full makes them -eligible without rescanning retained events. Eligible identity authors are the -owners and declared maintainers of accepted Full announcements plus accepted -root authors—not every author retained by the relay. A once-per-minute -reconciliation derives sources from this index. Remote queries contain at most 100 authors and -only one discovery batch may be in flight globally. Admission samples immediate +eligible without rescanning retained root events. A bounded recursive scan of +locally accepted root threads adds the authors of replies, reactions, zaps and +other descendants. Root provenance is carried through event-ID and address +references, so each participant is associated only with accepted threads in +which their events occur, including indirect descendants that expose only an +immediate parent. Eligible identity authors are +the owners and declared maintainers of accepted Full announcements plus +accepted root and descendant authors—not every author retained by the relay. A +once-per-minute reconciliation derives sources from this index. Remote queries +contain at most 100 authors and only one discovery batch may be in flight +globally. Admission samples immediate transient capacity; if historic work wins the small race before the permit is acquired, that single batch may wait but discovery can never build a waiter queue. Its SDK-owned REQ draws from the same pacer and subscription ledger as @@ -51,7 +62,7 @@ reconnection machinery, but are a distinct control-plane role: connecting to a user index or outbox does not start ordinary announcement, repository, or descendant sync against it. At most one new discovery source is dialled per maintenance pass, and an exclusively discovery connection retires after its -currently due author batches drain. A relay that independently becomes a +currently due identity or mailbox work drains. A relay that independently becomes a repository source is promoted to the ordinary lifecycle without opening a duplicate connection. Authors returned by a successful query refresh after 24 hours; missing authors and failed queries retry after five minutes. Configured user @@ -61,25 +72,64 @@ write/unmarked outboxes are then followed additively for newer replacements; new outboxes discovered by a replacement join the same bounded round until no unvisited source remains. Identity events pass through the ordinary write policy, persistence and broadcast path. Only a -stored, accepted kind `10002` can change inbox ownership. On startup, retained -relay lists for eligible authors rebuild inbox ownership from the local +stored, accepted kind `10002` can change mailbox ownership. On startup, retained +relay lists for eligible authors rebuild mailbox ownership from the local database before any network refresh, so serving established coverage does not depend on an external index remaining available. Remote discovery then refreshes that retained state on the normal cadence. -Inbox replacement/removal changes desired ownership immediately. Additions are -derived promptly. Removals prevent future rotating and historic work, but do -not eagerly CLOSE live descendants or abort shared in-flight batches: existing -coverage drains on its natural EOSE/CLOSED/disconnect or an ordinary later -consolidation rebuild. This avoids interrupting unrelated roots and subscription -churn merely because one author replaced a relay list. Root deletion follows -the existing GRASP-02 root-index lifecycle and is reconstructed from retained -accepted events on restart; this change does not add a second deletion graph. +Root-author read/unmarked inboxes still use ordinary GRASP-02 coverage. Wider +participant mailboxes use history-only workers. A maintenance pass starts at +most one due relay, while each relay may have at most one worker in flight. +Relays do not share a mailbox lane or terminal state: a slow or unavailable +relay cannot prevent another relay from progressing on a later pass. + +Each worker selects one stable-sorted, byte-bounded filter and delegates its +REQ lifecycle to the existing `RelayConnection::fetch_events` path. That path +owns per-relay request pacing, background priority, subscription-ledger +capacity, EOSE/CLOSED handling and a 30-second timeout. The worker reuses the +ordinary pagination state to continue through full historic pages until the +filter is exhausted or a page makes no new progress. It then passes events +through the normal write policy and persistence pipeline. No mailbox-specific +pending-batch kind, EOSE hook, close API, watchdog or cross-relay coordinator +is added. + +The probe covers accepted repository coordinates, root IDs, bounded recursive +descendant IDs, and descendant address coordinates using `a`/`A`/`q` and +`e`/`E`/`q`. A numeric in-memory cursor gives each filter group a turn. A +successful complete rotation refreshes after 24 hours; a failed group advances +the cursor after a five-minute delay and is retried on a later rotation. This +is deliberately best effort rather than a durable exactly-once schedule. A +relay already needed for ordinary repository sync shares its connection; an +exclusively control-plane connection retires when its identity and mailbox +work is idle. + +On restart, accepted roots and retained kind `10002` events rebuild participant +and mailbox ownership from LMDB. The in-memory group cursor is intentionally +not restored, so historic mailbox probing starts again from the first current +filter and remains safe to repeat. Inventory changes keep the cursor modulo the +new stable-sorted filter set; they do not coordinate with a worker on another +relay. + +The fixed-cardinality retained-state metric reports eligible authors, desired +history-probe relays, allocated filter cursors, and active per-relay workers. +This makes a stuck worker or unexpected inventory expansion visible without +putting peer URLs or public keys into metric labels. + +Mailbox replacement/removal changes desired ownership immediately. Root-author +inbox additions use ordinary coverage, while removal-only changes let existing +shared live subscriptions drain naturally. History-probe additions become due +promptly; removals prevent future groups while an already in-flight +`fetch_events` worker may finish naturally. Root deletion follows the existing +GRASP-02 root-index lifecycle and is reconstructed from retained accepted +events on restart; this change does not add a second deletion graph. ## Deliberately excluded -- the old recursive `SyncScope` graph and response-author fan-out; +- mailbox expansion for authors who have not produced locally accepted + repository-thread events; - maintainer mailbox expansion unrelated to an accepted root; -- identity storage for authors other than accepted roots; +- identity storage for authors outside accepted repositories and their + locally accepted threads; - per-user fallback configuration knobs; and -- a separate GRASP-03 subscription scheduler. +- permanent non-root participant-mailbox live subscriptions. diff --git a/docs/explanation/peer-controlled-state.md b/docs/explanation/peer-controlled-state.md index 27209a2..bf4b442 100644 --- a/docs/explanation/peer-controlled-state.md +++ b/docs/explanation/peer-controlled-state.md @@ -39,7 +39,7 @@ larger retained set; it is not an excuse to copy arbitrary wire input. | Purgatory event maps and sync queue | write policy; promotion/cleanup owns removal | Time/external: entries are keyed by admitted event/repository identity, normally expire at 30 minutes, and soft-expired announcements at 24 hours. Queue entries deduplicate by identifier and disappear on completion or event expiry. | purgatory counts, queue and Git-process metrics | | Per-domain Git throttle queues | incomplete purgatory fetch; throttle manager owns drain | External: one entry per purgatory identifier/domain, merged on repeat. Completion, URL exhaustion, or purgatory expiry removes useful work. Request history is time-windowed. | domain/fetch logs and Git-process metrics | | Dependency retry attempts and temporary relays | rejected/purgatory dependency discovery; maintenance owns expiry | Time/external: event IDs are pruned against current purgatory input; temporary relays have explicit deadlines. Repeats overwrite timestamps. | retained dependency gauges | -| NIP-65 discovery state/results | accepted root authors; discovery scheduler owns completion | External plus static work: author/source maps derive from accepted roots, one discovery query is in flight, and its result channel has capacity one. Missing results time out and clear in-flight state through result handling/disconnect refresh. | discovery logs and relay gauges | +| NIP-65 discovery and mailbox probes | accepted root and descendant authors; ordinary root-inbox coverage plus discovery/probe scheduler own completion | External/session work: root-author read inboxes remain ordinary derived sync targets, while participant maps derive from accepted root threads. Identity discovery is globally single-flight. Mailbox history has at most one ordinary `fetch_events` worker per accepted relay; each page consumes that relay's pacing and ledger capacity and has a 30-second timeout. Starts are paced one per maintenance pass, progress on different relays is independent, and exclusive connections retire when idle. Numeric filter cursors and due times are in memory and bounded by desired mailbox relays. | discovery/probe logs and relay gauges | | Deferred consolidation | capacity refusal; final batch/reset/disconnect owns removal | External and deduplicated by relay. Final batch completion processes the set directly; no self-addressed notification queue remains. | retained deferred-consolidation gauge | | Descendant rotations and auxiliary live coverage | accepted root coverage; EOSE/CLOSED/disconnect/daily reset own transition | Session/external: one state object per derived relay and at most one rotating request in flight per relay. Coverage IDs consume ledger slots. | retained rotation gauge and terminal logs | | Health and naughty-list entries | connection failures; health checker owns recovery/expiry | External/time: one entry per canonical target, ordinary failures back off, persistent entries expire after 12 hours. Metrics expose only three fixed categories. | health gauges and aggregate naughty metrics | diff --git a/src/nostr/builder.rs b/src/nostr/builder.rs index 16fd0ec..6f0accd 100644 --- a/src/nostr/builder.rs +++ b/src/nostr/builder.rs @@ -28,7 +28,8 @@ use crate::nostr::persistence::{EventPersistence, SaveContext}; use crate::nostr::policy::{ accepted_purgatory, duplicate, reject_error, reject_invalid, reject_restricted, AnnouncementPolicy, AnnouncementResult, IdentityAdmission, PolicyContext, PrEventPolicy, - ReferenceResult, RelatedEventPolicy, SharedProactiveRootAuthorIndex, StatePolicy, StateResult, + ReferenceResult, RelatedEventPolicy, SharedProactiveParticipantAuthorIndex, StatePolicy, + StateResult, }; use crate::nostr::SharedDatabase; use crate::purgatory::promotion_hooks::NostrPurgatoryPromotionHooks; @@ -84,7 +85,7 @@ pub struct Nip34WritePolicy { state_policy: StatePolicy, pr_event_policy: PrEventPolicy, related_event_policy: RelatedEventPolicy, - proactive_root_authors: SharedProactiveRootAuthorIndex, + proactive_participant_authors: SharedProactiveParticipantAuthorIndex, deletion: DeletionService, rejected_events_index: Arc>>>, } @@ -133,23 +134,24 @@ impl Nip34WritePolicy { state_policy: StatePolicy::new(ctx.clone()), pr_event_policy: PrEventPolicy::new(ctx.clone(), repo_init_locks), related_event_policy: RelatedEventPolicy::new(ctx.clone()), - proactive_root_authors: crate::nostr::policy::ProactiveRootAuthorIndex::shared(), + proactive_participant_authors: + crate::nostr::policy::ProactiveParticipantAuthorIndex::shared(), deletion: DeletionService::new(deletion_ctx), rejected_events_index: Arc::new(RwLock::new(None)), ctx, } } - pub fn proactive_root_authors(&self) -> SharedProactiveRootAuthorIndex { - self.proactive_root_authors.clone() + pub fn proactive_participant_authors(&self) -> SharedProactiveParticipantAuthorIndex { + self.proactive_participant_authors.clone() } async fn handle_proactive_identity(&self, event: &Event) -> WritePolicyResult { - match self.proactive_root_authors.admit(event).await { + match self.proactive_participant_authors.admit(event).await { IdentityAdmission::Accept => WritePolicyResult::Accept, - IdentityAdmission::IrrelevantAuthor => { - reject_restricted("Kind 0/10002 author must own an accepted repository root") - } + IdentityAdmission::IrrelevantAuthor => reject_restricted( + "Kind 0/10002 author must participate in an accepted repository thread", + ), IdentityAdmission::NotIdentity => { reject_invalid("Proactive identity policy received a non-identity event") } diff --git a/src/nostr/policy/identity.rs b/src/nostr/policy/identity.rs index 8be7e74..1b92d8c 100644 --- a/src/nostr/policy/identity.rs +++ b/src/nostr/policy/identity.rs @@ -1,4 +1,4 @@ -//! Narrow GRASP-03 identity admission for accepted root authors. +//! Narrow GRASP-03 identity admission for accepted repository participants. use std::collections::HashSet; use std::sync::Arc; @@ -6,7 +6,7 @@ use std::sync::Arc; use nostr_sdk::prelude::{Event, Kind, PublicKey}; use tokio::sync::RwLock; -pub type SharedProactiveRootAuthorIndex = Arc; +pub type SharedProactiveParticipantAuthorIndex = Arc; #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum IdentityAdmission { @@ -15,15 +15,15 @@ pub enum IdentityAdmission { NotIdentity, } -/// Derived admission set. Accepted repository roots remain authoritative; an -/// identity event can never add its own author here. +/// Derived admission set. Accepted repository threads remain authoritative; +/// an identity event can never add its own author here. #[derive(Debug, Default)] -pub struct ProactiveRootAuthorIndex { +pub struct ProactiveParticipantAuthorIndex { authors: RwLock>, } -impl ProactiveRootAuthorIndex { - pub fn shared() -> SharedProactiveRootAuthorIndex { +impl ProactiveParticipantAuthorIndex { + pub fn shared() -> SharedProactiveParticipantAuthorIndex { Arc::new(Self::default()) } @@ -52,7 +52,7 @@ mod tests { async fn identity_cannot_make_its_own_author_relevant() { let author = Keys::generate(); let other = Keys::generate(); - let index = ProactiveRootAuthorIndex::shared(); + let index = ProactiveParticipantAuthorIndex::shared(); let metadata = EventBuilder::new(Kind::Metadata, "{}") .finalize(&author) .unwrap(); diff --git a/src/nostr/policy/mod.rs b/src/nostr/policy/mod.rs index fa94dcf..c373cb7 100644 --- a/src/nostr/policy/mod.rs +++ b/src/nostr/policy/mod.rs @@ -18,7 +18,9 @@ pub(crate) use result::{ }; pub use announcement::{AnnouncementPolicy, AnnouncementResult}; -pub use identity::{IdentityAdmission, ProactiveRootAuthorIndex, SharedProactiveRootAuthorIndex}; +pub use identity::{ + IdentityAdmission, ProactiveParticipantAuthorIndex, SharedProactiveParticipantAuthorIndex, +}; pub use pr_event::PrEventPolicy; pub use related::{ReferenceResult, RelatedEventPolicy}; pub use state::{StatePolicy, StateResult}; diff --git a/src/sync/discovery.rs b/src/sync/discovery.rs index 4d2b267..b4d9183 100644 --- a/src/sync/discovery.rs +++ b/src/sync/discovery.rs @@ -175,6 +175,25 @@ pub fn accepted_repository_authors<'a>( authors } +/// Associate each accepted root author with their own roots and each accepted +/// descendant author with the root provenance derived for their events. +pub fn mailbox_author_roots( + roots: &[AcceptedRoot], + participant_roots: impl IntoIterator)>, +) -> HashMap> { + let mut author_roots: HashMap> = HashMap::new(); + for root in roots { + author_roots.entry(root.author).or_default().insert(root.id); + } + for (author, accepted_roots) in participant_roots { + author_roots + .entry(author) + .or_default() + .extend(accepted_roots); + } + author_roots +} + pub fn merge_inbox_roots( targets: &mut HashMap, inbox_roots: &HashMap>, @@ -220,6 +239,15 @@ pub fn build_inbox_root_overlay( overlay } +/// Relays on which an author may receive or publish conversation events. +/// Unmarked entries serve both roles under NIP-65. +pub fn mailbox_relays(event: &Event) -> HashSet { + inbox_relays(event) + .into_iter() + .chain(outbox_relays(event)) + .collect() +} + #[cfg(test)] mod tests { use super::*; @@ -312,6 +340,29 @@ mod tests { ); } + #[test] + fn relay_list_maps_read_write_and_unmarked_participant_mailboxes() { + let keys = Keys::generate(); + let relay_list = event( + &keys, + Kind::RelayList, + vec![ + Tag::custom("r", ["wss://read.example", "read"]), + Tag::custom("r", ["wss://write.example", "write"]), + Tag::custom("r", ["wss://both.example"]), + ], + 1, + ); + assert_eq!( + mailbox_relays(&relay_list), + HashSet::from([ + "wss://read.example".to_string(), + "wss://write.example".to_string(), + "wss://both.example".to_string(), + ]) + ); + } + #[test] fn identity_authors_are_bounded_to_full_repo_owners_and_maintainers() { let owner = Keys::generate(); @@ -356,6 +407,40 @@ mod tests { ); } + #[test] + fn descendant_authors_receive_only_their_derived_root_provenance() { + let first_root_author = Keys::generate().public_key(); + let second_root_author = Keys::generate().public_key(); + let participant = Keys::generate().public_key(); + let first_root = EventId::from_byte_array([12; 32]); + let second_root = EventId::from_byte_array([13; 32]); + let roots = vec![ + AcceptedRoot { + id: first_root, + author: first_root_author, + repository: "30617:owner:first".to_string(), + }, + AcceptedRoot { + id: second_root, + author: second_root_author, + repository: "30617:owner:second".to_string(), + }, + ]; + + let author_roots = + mailbox_author_roots(&roots, [(participant, HashSet::from([second_root]))]); + + assert_eq!( + author_roots[&first_root_author], + HashSet::from([first_root]) + ); + assert_eq!( + author_roots[&second_root_author], + HashSet::from([second_root]) + ); + assert_eq!(author_roots[&participant], HashSet::from([second_root])); + } + #[test] fn latest_relay_list_uses_nip01_replacement_order() { let keys = Keys::generate(); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 18a6f91..a54a86e 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -68,6 +68,123 @@ const DESCENDANT_FALLBACK_OVERLAP_SECS: u64 = 15 * 60; /// Maximum number of locally known parent/child generations expanded into a /// relay's related-event query frontier. const MAX_DESCENDANT_FRONTIER_DEPTH: usize = 8; + +fn mailbox_probe_refresh_interval() -> Duration { + if std::env::var("NGIT_TEST").as_deref() == Ok("1") { + Duration::from_secs(2) + } else { + Duration::from_secs(24 * 60 * 60) + } +} + +fn mailbox_probe_retry_interval() -> Duration { + if std::env::var("NGIT_TEST").as_deref() == Ok("1") { + Duration::from_millis(500) + } else { + Duration::from_secs(5 * 60) + } +} + +fn due_mailbox_relays( + mailbox_roots: &HashMap>, + next_at: &HashMap, + now: Instant, +) -> Vec { + let mut relays: Vec = mailbox_roots + .keys() + .filter(|relay| next_at.get(*relay).is_none_or(|due| *due <= now)) + .cloned() + .collect(); + relays.sort_by(|left, right| { + next_at + .get(left) + .copied() + .unwrap_or(now) + .cmp(&next_at.get(right).copied().unwrap_or(now)) + .then_with(|| mailbox_roots[right].len().cmp(&mailbox_roots[left].len())) + .then_with(|| left.cmp(right)) + }); + relays +} + +fn mailbox_probe_connection_ready( + connection_status: Option, + socket_connected: bool, +) -> bool { + socket_connected && connection_status.is_some_and(|status| status.is_live_sync_active()) +} + +fn select_due_mailbox_relay( + due_relays: &[String], + active_relays: &HashSet, +) -> Option { + due_relays + .iter() + .find(|relay| active_relays.contains(*relay)) + .or_else(|| due_relays.first()) + .cloned() +} + +fn mailbox_probe_completion( + filter_index: usize, + filter_count: usize, + succeeded: bool, +) -> (usize, bool, Duration) { + let next_filter = (filter_index + 1) % filter_count.max(1); + let completed_cycle = next_filter == 0; + let next_probe_in = if succeeded && completed_cycle { + mailbox_probe_refresh_interval() + } else if succeeded { + Duration::ZERO + } else { + mailbox_probe_retry_interval() + }; + (next_filter, completed_cycle, next_probe_in) +} + +async fn fetch_mailbox_filter( + connection: RelayConnection, + filter: Filter, +) -> Result, String> { + let mut pagination = PaginationState::new(vec![filter]); + let mut session = RelayPaginationSession::default(); + let mut events = HashMap::::new(); + + loop { + let filter = pagination + .filters() + .into_iter() + .next() + .expect("mailbox pagination always owns one filter"); + let page = connection + .fetch_events(filter, Duration::from_secs(30)) + .await?; + let previous_count = events.len(); + for event in page { + pagination.record_event(&event); + events.insert(event.id, event); + } + // Inclusive `until` cursors can repeat a relay's oldest timestamp. + // Stop when a page contributes nothing instead of retaining the only + // mailbox worker forever on a deterministic repeated page. + if events.len() == previous_count { + break; + } + let Some(next) = pagination.next_page(&mut session) else { + break; + }; + pagination = next; + } + + let mut events: Vec<_> = events.into_values().collect(); + events.sort_by(|left, right| { + left.created_at + .cmp(&right.created_at) + .then_with(|| left.id.cmp(&right.id)) + }); + Ok(events) +} + fn should_use_semantic_fallback(requested_count: usize, received_count: usize) -> bool { requested_count >= SEMANTIC_FALLBACK_MIN_REQUESTED_EVENTS && received_count.saturating_mul(100) @@ -860,6 +977,9 @@ pub enum PendingBatchPurpose { struct DescendantFrontier { event_ids: HashSet, coordinates: HashSet, + /// Accepted participant authors and the exact repository roots reached by + /// their retained direct or recursive events. + author_roots: HashMap>, } #[derive(Debug, Default)] @@ -1006,6 +1126,32 @@ fn descendant_event_coordinate(event: &Event) -> Option { )) } +fn descendant_reference_roots( + event: &Event, + root_events: &HashSet, + event_root_owners: &HashMap>, + coordinate_root_owners: &HashMap>, +) -> HashSet { + let (coordinates, event_ids) = + crate::nostr::policy::RelatedEventPolicy::extract_reference_tags(event); + let mut roots: HashSet = event_ids + .iter() + .filter(|event_id| root_events.contains(event_id)) + .copied() + .collect(); + for event_id in event_ids { + if let Some(owners) = event_root_owners.get(&event_id) { + roots.extend(owners.iter().copied()); + } + } + for coordinate in coordinates { + if let Some(owners) = coordinate_root_owners.get(&coordinate) { + roots.extend(owners.iter().copied()); + } + } + roots +} + /// Derive a bounded transitive frontier from related events already accepted /// into the local database. /// @@ -1039,8 +1185,13 @@ async fn recursive_descendant_frontier_with_limits( } // Every event directly tagging a root owns an independent recursive - // subtree. The direct event itself remains complete core coverage. - let mut direct_events = HashMap::::new(); + // subtree. This keeps a popular tangential branch from consuming the + // compatibility coverage of unrelated issues, patches, or repositories. + // Retain exact root provenance at the same time so participant mailbox + // probes cannot leak one repository's roots into another author's scope. + let mut direct_events = HashMap::)>::new(); + let mut event_root_owners = HashMap::>::new(); + let mut coordinate_root_owners = HashMap::>::new(); for filter in filters::tagged_one_of_our_root_event_filters(root_events, None) { match database.query(filter).await { Ok(events) => { @@ -1048,9 +1199,21 @@ async fn recursive_descendant_frontier_with_limits( .iter() .filter(|event| !root_events.contains(&event.id)) { + let roots = descendant_reference_roots( + event, + root_events, + &event_root_owners, + &coordinate_root_owners, + ); + if roots.is_empty() { + continue; + } direct_events .entry(event.id) - .or_insert_with(|| event.clone()); + .and_modify(|(_, known_roots)| { + known_roots.extend(roots.iter().copied()); + }) + .or_insert_with(|| (event.clone(), roots)); } } Err(error) => { @@ -1065,7 +1228,7 @@ async fn recursive_descendant_frontier_with_limits( } let mut direct_events: Vec<_> = direct_events.into_iter().collect(); - direct_events.sort_by(|(left_id, left), (right_id, right)| { + direct_events.sort_by(|(left_id, (left, _)), (right_id, (right, _))| { left.created_at .cmp(&right.created_at) .then_with(|| left_id.cmp(right_id)) @@ -1076,9 +1239,15 @@ async fn recursive_descendant_frontier_with_limits( let mut coordinate_branches = HashMap::>::new(); let mut event_seeds = HashMap::>::new(); let mut coordinate_seeds = HashMap::>::new(); - for (event_id, event) in direct_events { + for (event_id, (event, roots)) in direct_events { branch_member_counts.insert(event_id, 0); event_branches.entry(event_id).or_default().insert(event_id); + frontier + .author_roots + .entry(event.pubkey) + .or_default() + .extend(roots.iter().copied()); + event_root_owners.insert(event_id, roots.clone()); event_seeds.entry(event_id).or_default().insert(event_id); if let Some(coordinate) = descendant_event_coordinate(&event) { coordinate_branches @@ -1086,9 +1255,13 @@ async fn recursive_descendant_frontier_with_limits( .or_default() .insert(event_id); coordinate_seeds - .entry(coordinate) + .entry(coordinate.clone()) .or_default() .insert(event_id); + coordinate_root_owners + .entry(coordinate) + .or_default() + .extend(roots); } } @@ -1099,7 +1272,7 @@ async fn recursive_descendant_frontier_with_limits( } for depth in 2..=max_depth { - let mut layer_events = HashMap::::new(); + let mut layer_events = HashMap::)>::new(); let event_seed_ids: HashSet<_> = event_seeds.keys().copied().collect(); let coordinate_seed_ids: HashSet<_> = coordinate_seeds.keys().cloned().collect(); let mut layer_filters = @@ -1113,13 +1286,24 @@ async fn recursive_descendant_frontier_with_limits( match database.query(filter).await { Ok(events) => { for event in events.iter() { - if root_events.contains(&event.id) - || direct_event_ids.contains(&event.id) - || layer_events.contains_key(&event.id) - { + if root_events.contains(&event.id) || direct_event_ids.contains(&event.id) { continue; } - layer_events.insert(event.id, event.clone()); + let roots = descendant_reference_roots( + event, + root_events, + &event_root_owners, + &coordinate_root_owners, + ); + if roots.is_empty() { + continue; + } + layer_events + .entry(event.id) + .and_modify(|(_, known_roots)| { + known_roots.extend(roots.iter().copied()); + }) + .or_insert_with(|| (event.clone(), roots)); } } Err(error) => { @@ -1137,7 +1321,7 @@ async fn recursive_descendant_frontier_with_limits( } let mut layer_events: Vec<_> = layer_events.into_iter().collect(); - layer_events.sort_by(|(left_id, left), (right_id, right)| { + layer_events.sort_by(|(left_id, (left, _)), (right_id, (right, _))| { left.created_at .cmp(&right.created_at) .then_with(|| left_id.cmp(right_id)) @@ -1145,7 +1329,7 @@ async fn recursive_descendant_frontier_with_limits( let mut next_event_seeds = HashMap::>::new(); let mut next_coordinate_seeds = HashMap::>::new(); - for (event_id, event) in layer_events { + for (event_id, (event, roots)) in layer_events { let (addressable_refs, event_refs) = crate::nostr::policy::RelatedEventPolicy::extract_reference_tags(&event); let mut inherited_branches = HashSet::new(); @@ -1174,15 +1358,29 @@ async fn recursive_descendant_frontier_with_limits( if admitted_branches.is_empty() { continue; } + frontier + .author_roots + .entry(event.pubkey) + .or_default() + .extend(roots.iter().copied()); + event_root_owners + .entry(event_id) + .or_default() + .extend(roots.iter().copied()); event_branches .entry(event_id) .or_default() .extend(admitted_branches.iter().copied()); - if let Some(coordinate) = descendant_event_coordinate(&event) { + let coordinate = descendant_event_coordinate(&event); + if let Some(coordinate) = &coordinate { coordinate_branches - .entry(coordinate) + .entry(coordinate.clone()) .or_default() .extend(admitted_branches.iter().copied()); + coordinate_root_owners + .entry(coordinate.clone()) + .or_default() + .extend(roots.iter().copied()); } let expandable_branches: HashSet<_> = admitted_branches @@ -1193,7 +1391,7 @@ async fn recursive_descendant_frontier_with_limits( continue; } next_event_seeds.insert(event_id, expandable_branches.clone()); - if let Some(coordinate) = descendant_event_coordinate(&event) { + if let Some(coordinate) = coordinate { next_coordinate_seeds.insert(coordinate, expandable_branches); } } @@ -1251,6 +1449,22 @@ fn descendant_frontier_filters( filters } +fn mailbox_probe_filters( + repositories: &HashSet, + root_events: &HashSet, + frontier: &DescendantFrontier, +) -> Vec { + let coordinate_values: HashSet<_> = + repositories.union(&frontier.coordinates).cloned().collect(); + let event_values: HashSet<_> = root_events.union(&frontier.event_ids).copied().collect(); + let mut filters = filters::tagged_one_of_our_repo_event_filters(&coordinate_values, None); + filters.extend(filters::tagged_one_of_our_root_event_filters( + &event_values, + None, + )); + filters +} + fn packed_historic_filters(items: &PendingItems, frontier: &DescendantFrontier) -> Vec { let all_repos: HashSet<_> = items .repos @@ -1609,23 +1823,68 @@ struct Nip65DiscoveryResult { outcome: Result, String>, } +#[derive(Debug)] +struct MailboxProbeResult { + source_relay: String, + filter_index: usize, + filter_count: usize, + outcome: Result, String>, +} + #[derive(Debug, Default)] struct Nip65DiscoveryState { + /// Roots authored by each accepted root author. These preserve the + /// existing GRASP-03 live inbox overlay; non-root participants remain on + /// the paced history-only mailbox path below. + root_author_roots: HashMap>, author_roots: HashMap>, + root_repositories: HashMap, eligible_authors: HashSet, author_sources: HashMap>, relay_lists: HashMap, author_inboxes: HashMap>, + author_mailboxes: HashMap>, /// Eligible authors for whom at least one successful index query returned /// no accepted relay list. They use the operator's bounded fallback set /// until an accepted kind 10002 arrives. fallback_authors: HashSet, inbox_roots: HashMap>, + mailbox_roots: HashMap>, + mailbox_probe_next_at: HashMap, + mailbox_probe_next_filter: HashMap, + /// At most one ordinary fetch runs per mailbox relay. Different relays + /// remain independent and use their own connection capacity. + mailbox_probes_in_flight: HashSet, in_flight: HashSet<(String, PublicKey)>, next_attempt_at: HashMap<(String, PublicKey), Instant>, next_inventory_at: Option, } +impl Nip65DiscoveryState { + fn install_mailbox_overlay( + &mut self, + new_overlay: HashMap>, + now: Instant, + ) { + let old_overlay = std::mem::replace(&mut self.mailbox_roots, new_overlay); + if old_overlay == self.mailbox_roots { + return; + } + self.mailbox_probe_next_at + .retain(|relay, _| self.mailbox_roots.contains_key(relay)); + self.mailbox_probe_next_filter + .retain(|relay, _| self.mailbox_roots.contains_key(relay)); + for (relay, roots) in &self.mailbox_roots { + if old_overlay.get(relay) != Some(roots) { + self.mailbox_probe_next_at.insert(relay.clone(), now); + self.mailbox_probe_next_filter + .entry(relay.clone()) + .or_default(); + } + } + } +} + /// Quick reconnect window in seconds (15 minutes) const QUICK_RECONNECT_WINDOW_SECS: u64 = 15 * 60; @@ -2050,12 +2309,17 @@ async fn run_health_and_metrics_checker( // persistent set. manager.sync_due_byte_limited_relay().await; - // 5. Discover one bounded batch of root-author NIP-65 inboxes. + // 5. Discover one bounded batch of participant NIP-65 mailboxes. // The query uses an ordinary transient ledger slot and yields // to historic work when no slot is immediately available. manager.schedule_nip65_discovery().await; - // 6. Check for naughty list expiration + // 6. Start at most one participant mailbox history filter. + // Mailbox coverage is history-only: it does not turn every + // participant relay into a permanent live source. + manager.schedule_mailbox_probe().await; + + // 7. Check for naughty list expiration if let Some(naughty_list) = manager.health_tracker.naughty_list() { let recovered = naughty_list.expire_old_entries(); for url in recovered { @@ -2066,7 +2330,7 @@ async fn run_health_and_metrics_checker( } } - // 7. Update metrics with current health states and naughty list + // 8. Update metrics with current health states and naughty list if let Some(ref metrics) = manager.metrics { // Get all tracked relay URLs let relay_urls: Vec = { @@ -2107,6 +2371,22 @@ async fn run_health_and_metrics_checker( ("related_dependency_events", manager.rejected_events_index.related_len()), ("deferred_consolidations", manager.deferred_consolidations.relays.len()), ("descendant_rotations", manager.descendant_sync_rotations.len()), + ( + "nip65_eligible_authors", + manager.nip65_discovery.eligible_authors.len(), + ), + ( + "mailbox_probe_relays", + manager.nip65_discovery.mailbox_roots.len(), + ), + ( + "mailbox_probe_cursors", + manager.nip65_discovery.mailbox_probe_next_filter.len(), + ), + ( + "mailbox_probe_active", + manager.nip65_discovery.mailbox_probes_in_flight.len(), + ), ] { metrics.set_retained_state(class, count); } @@ -2143,7 +2423,7 @@ pub struct SyncManager { /// What we want to sync (source of truth) repo_sync_index: RepoSyncIndex, root_candidate_index: RootCandidateIndex, - proactive_root_authors: crate::nostr::policy::SharedProactiveRootAuthorIndex, + proactive_participant_authors: crate::nostr::policy::SharedProactiveParticipantAuthorIndex, /// What we've confirmed syncing + connection state relay_sync_index: RelaySyncIndex, /// In-flight subscription batches @@ -2205,8 +2485,7 @@ pub struct SyncManager { /// Round-robin cursor so only one relay performs descendant live-mode /// derivation and admission work per maintenance tick. descendant_relay_cursor: usize, - /// Narrow GRASP-03 control-plane state. Event transport remains in the - /// ordinary GRASP-02 relay targets after inbox roots are overlaid. + /// Narrow GRASP-03 identity discovery and paced mailbox-probe state. nip65_discovery: Nip65DiscoveryState, /// Channel for disconnect notifications (set during run) disconnect_tx: Option>, @@ -2217,6 +2496,7 @@ pub struct SyncManager { /// Returns connection outcomes to the sync actor for serialized state changes. connect_attempt_result_tx: Option>, nip65_discovery_result_tx: Option>, + mailbox_probe_result_tx: Option>, /// Channel for broadcasting shutdown signal to all background tasks shutdown_tx: Option>, /// Prometheus metrics for sync operations (None if metrics disabled) @@ -2279,7 +2559,7 @@ impl SyncManager { } } - let proactive_root_authors = write_policy.proactive_root_authors(); + let proactive_participant_authors = write_policy.proactive_participant_authors(); Self { bootstrap_relay_url, service_domain, @@ -2290,7 +2570,7 @@ impl SyncManager { config: config.clone(), repo_sync_index: Arc::new(RwLock::new(HashMap::new())), root_candidate_index: Arc::new(RwLock::new(HashMap::new())), - proactive_root_authors, + proactive_participant_authors, relay_sync_index: Arc::new(RwLock::new(HashMap::new())), pending_sync_index: Arc::new(RwLock::new(HashMap::new())), rejected_events_index, @@ -2321,6 +2601,7 @@ impl SyncManager { subscription_closed_tx: None, connect_attempt_result_tx: None, nip65_discovery_result_tx: None, + mailbox_probe_result_tx: None, shutdown_tx: None, metrics: sync_metrics, } @@ -3736,6 +4017,8 @@ impl SyncManager { let (nip65_discovery_result_tx, mut nip65_discovery_result_rx) = mpsc::channel::(1); + let (mailbox_probe_result_tx, mut mailbox_probe_result_rx) = + mpsc::channel::(1); // 4b. Create shutdown broadcast channel for graceful shutdown let (shutdown_tx, _shutdown_rx) = broadcast::channel(1); @@ -3758,6 +4041,7 @@ impl SyncManager { self.subscription_closed_tx = Some(subscription_closed_tx); self.connect_attempt_result_tx = Some(connect_attempt_result_tx); self.nip65_discovery_result_tx = Some(nip65_discovery_result_tx); + self.mailbox_probe_result_tx = Some(mailbox_probe_result_tx); self.shutdown_tx = Some(shutdown_tx.clone()); // 6. Connect to bootstrap relay if configured @@ -3860,6 +4144,12 @@ impl SyncManager { ).await; } } + result = mailbox_probe_result_rx.recv() => { + if let Some(result) = result { + let mut manager = sync_manager.lock().await; + manager.handle_mailbox_probe_result(result).await; + } + } result = connect_attempt_result_rx.recv() => { match result { Some(result) => { @@ -3886,9 +4176,15 @@ impl SyncManager { // state instead of trusting a queued full-index // snapshot that may already be stale. let mut manager = sync_manager.lock().await; + // A dirty relay can contain a newly accepted root + // or participant. Refresh the local inventory now; + // the per-author retry deadlines still bound remote + // NIP-65 queries when only repository state changed. + manager.nip65_discovery.next_inventory_at = None; manager .recompute_new_sync_filters_for_relay(&add_filters.relay_url) .await; + manager.schedule_nip65_discovery().await; } None => break, } @@ -5157,7 +5453,7 @@ impl SyncManager { async fn schedule_nip65_discovery(&mut self) { // GRASP-03 is an optional overlay on the always-running GRASP-02 // manager. When disabled, do not inventory authors, retain identity - // authority, open discovery connections, or add inbox roots. + // authority, open discovery connections, or add mailbox roots. if !self.config.sync_plus_enabled { return; } @@ -5168,32 +5464,48 @@ impl SyncManager { .next_inventory_at .is_none_or(|due| due <= now) { - let candidates = self.root_candidate_index.read().await; - let repo_index = self.repo_sync_index.read().await; - let roots = discovery::accepted_root_candidates(candidates.values(), &repo_index); - drop(candidates); - let mut author_roots: HashMap> = HashMap::new(); - for root in roots { - author_roots.entry(root.author).or_default().insert(root.id); - } - let announcements = self .database .query(Filter::new().kind(Kind::GitRepoAnnouncement)) .await .unwrap_or_default(); + let candidates = self.root_candidate_index.read().await; + let repo_index = self.repo_sync_index.read().await; + let roots = discovery::accepted_root_candidates(candidates.values(), &repo_index); let mut eligible_authors = discovery::accepted_repository_authors(announcements.iter(), &repo_index); - eligible_authors.extend(author_roots.keys().copied()); + drop(candidates); drop(repo_index); + let root_ids: HashSet = roots.iter().map(|root| root.id).collect(); + self.nip65_discovery.root_repositories = roots + .iter() + .map(|root| (root.id, root.repository.clone())) + .collect(); + self.nip65_discovery.root_author_roots = + discovery::mailbox_author_roots(&roots, std::iter::empty()); + + // A root may be absent from the author's write relays while a + // participant's reaction, zap or reply is present there. Carry + // root provenance through indirect event/address descendants so + // each mailbox probe stays scoped to accepted threads in which + // that author actually participated. + let frontier = recursive_descendant_frontier( + &self.database, + &root_ids, + self.config.sync_recursive_descendant_limit, + ) + .await; + let participant_author_count = frontier.author_roots.len(); + let author_roots = discovery::mailbox_author_roots(&roots, frontier.author_roots); + eligible_authors.extend(author_roots.keys().copied()); self.nip65_discovery.author_roots = author_roots; self.nip65_discovery.eligible_authors = eligible_authors; - self.proactive_root_authors + self.proactive_participant_authors .replace(self.nip65_discovery.eligible_authors.clone()) .await; let current_authors = self.nip65_discovery.eligible_authors.clone(); - // Rebuild inbox ownership from already accepted local state before + // Rebuild mailbox ownership from already accepted local state before // attempting a network refresh. GRASP-03 retains kind 10002 so a // restart must not make conversation coverage depend on the index // source still being reachable. @@ -5223,6 +5535,12 @@ impl SyncManager { .iter() .map(|(author, event)| (*author, discovery::inbox_relays(event))) .collect(); + self.nip65_discovery.author_mailboxes = self + .nip65_discovery + .relay_lists + .iter() + .map(|(author, event)| (*author, discovery::mailbox_relays(event))) + .collect(); } let mut index_relays: HashSet = self @@ -5264,6 +5582,9 @@ impl SyncManager { self.nip65_discovery .author_inboxes .retain(|author, _| current_authors.contains(author)); + self.nip65_discovery + .author_mailboxes + .retain(|author, _| current_authors.contains(author)); self.nip65_discovery.fallback_authors.retain(|author| { current_authors.contains(author) && !self.nip65_discovery.relay_lists.contains_key(author) @@ -5277,13 +5598,29 @@ impl SyncManager { .is_some_and(|sources| sources.contains(relay)) }); self.nip65_discovery.next_inventory_at = Some(now + nip65_inventory_interval()); - let overlay = discovery::build_inbox_root_overlay( - &self.nip65_discovery.author_roots, + let fallback_relays = self.configured_nip65_fallback_relays(); + let inbox_overlay = discovery::build_inbox_root_overlay( + &self.nip65_discovery.root_author_roots, &self.nip65_discovery.author_inboxes, &self.nip65_discovery.fallback_authors, - &self.configured_nip65_fallback_relays(), + &fallback_relays, + ); + let mailbox_overlay = discovery::build_inbox_root_overlay( + &self.nip65_discovery.author_roots, + &self.nip65_discovery.author_mailboxes, + &self.nip65_discovery.fallback_authors, + &fallback_relays, + ); + self.install_nip65_inbox_overlay(inbox_overlay).await; + self.install_nip65_mailbox_overlay(mailbox_overlay); + tracing::info!( + root_count = root_ids.len(), + participant_author_count, + eligible_author_count = self.nip65_discovery.eligible_authors.len(), + mailbox_relay_count = self.nip65_discovery.mailbox_roots.len(), + fallback_authors = self.nip65_discovery.fallback_authors.len(), + "Reconciled proactive participant mailbox inventory" ); - self.install_nip65_overlay(overlay).await; } // Discovery is deliberately single-flight. The immediate-capacity @@ -5370,7 +5707,10 @@ impl SyncManager { } } - async fn install_nip65_overlay(&mut self, new_overlay: HashMap>) { + async fn install_nip65_inbox_overlay( + &mut self, + new_overlay: HashMap>, + ) { let old_overlay = std::mem::replace(&mut self.nip65_discovery.inbox_roots, new_overlay); if old_overlay == self.nip65_discovery.inbox_roots { return; @@ -5379,14 +5719,213 @@ impl SyncManager { dirty_relays.extend(self.nip65_discovery.inbox_roots.keys().cloned()); for relay in dirty_relays { - // Desired ownership changes immediately, but existing descendant - // work retires only at its natural EOSE/CLOSED/disconnect or a - // later ordinary consolidation. Closing shared batches here can - // interrupt unrelated roots and churn otherwise healthy live REQs. + // Preserve the established root-author inbox behavior: additions + // enter ordinary live/rotating descendant coverage, while + // removal-only changes drain through that coverage's normal + // lifecycle instead of churning shared subscriptions. self.recompute_new_sync_filters_for_relay(&relay).await; } } + fn install_nip65_mailbox_overlay(&mut self, new_overlay: HashMap>) { + self.nip65_discovery + .install_mailbox_overlay(new_overlay, Instant::now()); + } + + async fn schedule_mailbox_probe(&mut self) { + if !self.config.sync_plus_enabled { + return; + } + + let now = Instant::now(); + let due_relays: Vec = due_mailbox_relays( + &self.nip65_discovery.mailbox_roots, + &self.nip65_discovery.mailbox_probe_next_at, + now, + ) + .into_iter() + .filter(|relay| { + !self + .nip65_discovery + .mailbox_probes_in_flight + .contains(relay) + }) + .collect(); + let active_relays: HashSet = self + .relay_sync_index + .read() + .await + .iter() + .filter(|(_, state)| state.connection_status.is_live_sync_active()) + .map(|(relay, _)| relay.clone()) + .collect(); + let Some(relay) = select_due_mailbox_relay(&due_relays, &active_relays) else { + return; + }; + if is_own_sync_target(&relay, &self.service_domain) + || self.rejected_relay_targets.contains(&relay) + { + self.nip65_discovery + .mailbox_probe_next_at + .insert(relay, now + mailbox_probe_refresh_interval()); + return; + } + + if !self.connections.contains_key(&relay) { + if self.register_relay(relay.clone(), false, true).await { + self.schedule_connect_relay(&relay).await; + } + self.defer_mailbox_relay(&relay); + return; + } + + let Some(connection) = self.connections.get(&relay).cloned() else { + self.defer_mailbox_relay(&relay); + return; + }; + let socket_connected = connection.is_connected().await; + let connection_status = self + .relay_sync_index + .read() + .await + .get(&relay) + .map(|state| state.connection_status); + if !mailbox_probe_connection_ready(connection_status, socket_connected) { + if !socket_connected && self.health_tracker.should_attempt_connection(&relay) { + self.schedule_connect_relay(&relay).await; + } else if !socket_connected { + self.retire_idle_nip65_discovery_source(&relay).await; + } + self.defer_mailbox_relay(&relay); + return; + } + + let roots = self.nip65_discovery.mailbox_roots[&relay].clone(); + let repositories: HashSet = roots + .iter() + .filter_map(|root| self.nip65_discovery.root_repositories.get(root).cloned()) + .collect(); + let frontier = recursive_descendant_frontier( + &self.database, + &roots, + self.config.sync_recursive_descendant_limit, + ) + .await; + let mut filters = mailbox_probe_filters(&repositories, &roots, &frontier); + filters.sort_by_key(Filter::as_json); + let Some(filter_count) = (!filters.is_empty()).then_some(filters.len()) else { + self.nip65_discovery + .mailbox_probe_next_at + .insert(relay.clone(), now + mailbox_probe_refresh_interval()); + self.retire_idle_nip65_discovery_source(&relay).await; + return; + }; + let filter_index = self + .nip65_discovery + .mailbox_probe_next_filter + .get(&relay) + .copied() + .unwrap_or_default() + % filter_count; + let filter = filters.swap_remove(filter_index); + let Some(result_tx) = self.mailbox_probe_result_tx.clone() else { + self.defer_mailbox_relay(&relay); + return; + }; + self.nip65_discovery + .mailbox_probes_in_flight + .insert(relay.clone()); + let source_relay = relay.clone(); + tokio::spawn(async move { + let outcome = fetch_mailbox_filter(connection, filter).await; + let _ = result_tx + .send(MailboxProbeResult { + source_relay, + filter_index, + filter_count, + outcome, + }) + .await; + }); + tracing::info!( + relay = %relay, + filter_index, + filter_count, + "Started bounded participant mailbox fetch" + ); + } + + fn defer_mailbox_relay(&mut self, relay: &str) { + if self.nip65_discovery.mailbox_roots.contains_key(relay) { + self.nip65_discovery.mailbox_probe_next_at.insert( + relay.to_string(), + Instant::now() + mailbox_probe_retry_interval(), + ); + } + } + + async fn handle_mailbox_probe_result(&mut self, result: MailboxProbeResult) { + if !self + .nip65_discovery + .mailbox_probes_in_flight + .remove(&result.source_relay) + { + tracing::debug!( + relay = %result.source_relay, + "Ignoring stale participant mailbox result" + ); + return; + } + + let (events, succeeded) = match result.outcome { + Ok(events) => (events, true), + Err(error) => { + tracing::debug!(relay = %result.source_relay, %error, "Participant mailbox fetch failed"); + (Vec::new(), false) + } + }; + let event_count = events.len(); + for event in events { + let _ = Self::process_event_static( + &event, + &result.source_relay, + &self.database, + &self.write_policy, + &self.local_relay, + &self.rejected_events_index, + crate::nostr::persistence::SaveContext::RelaySync, + ) + .await; + } + + let (next_filter, completed_cycle, next_probe_in) = + mailbox_probe_completion(result.filter_index, result.filter_count, succeeded); + self.nip65_discovery + .mailbox_probe_next_filter + .insert(result.source_relay.clone(), next_filter); + if self + .nip65_discovery + .mailbox_roots + .contains_key(&result.source_relay) + { + self.nip65_discovery + .mailbox_probe_next_at + .insert(result.source_relay.clone(), Instant::now() + next_probe_in); + } + tracing::info!( + relay = %result.source_relay, + filter_index = result.filter_index, + filter_count = result.filter_count, + event_count, + succeeded, + completed_cycle, + next_probe_in_secs = next_probe_in.as_secs_f64(), + "Participant mailbox fetch reached terminal state" + ); + self.retire_idle_nip65_discovery_source(&result.source_relay) + .await; + } + async fn handle_nip65_discovery_result(&mut self, result: Nip65DiscoveryResult) { for author in &result.authors { self.nip65_discovery @@ -5469,6 +6008,9 @@ impl SyncManager { self.nip65_discovery .author_inboxes .insert(author, discovery::inbox_relays(&candidate)); + self.nip65_discovery + .author_mailboxes + .insert(author, discovery::mailbox_relays(&candidate)); let mut index_relays: HashSet = self .config .parse_user_index_relays() @@ -5503,19 +6045,27 @@ impl SyncManager { self.schedule_nip65_discovery().await; return; } - let overlay = discovery::build_inbox_root_overlay( - &self.nip65_discovery.author_roots, + let fallback_relays = self.configured_nip65_fallback_relays(); + let inbox_overlay = discovery::build_inbox_root_overlay( + &self.nip65_discovery.root_author_roots, &self.nip65_discovery.author_inboxes, &self.nip65_discovery.fallback_authors, - &self.configured_nip65_fallback_relays(), + &fallback_relays, ); - self.install_nip65_overlay(overlay).await; + let mailbox_overlay = discovery::build_inbox_root_overlay( + &self.nip65_discovery.author_roots, + &self.nip65_discovery.author_mailboxes, + &self.nip65_discovery.fallback_authors, + &fallback_relays, + ); + self.install_nip65_inbox_overlay(inbox_overlay).await; + self.install_nip65_mailbox_overlay(mailbox_overlay); tracing::info!( source = %result.source_relay, authors = result.authors.len(), - inbox_relays = self.nip65_discovery.inbox_roots.len(), + mailbox_relays = self.nip65_discovery.mailbox_roots.len(), fallback_authors = self.nip65_discovery.fallback_authors.len(), - "Updated proactive inbox coverage from NIP-65" + "Updated proactive participant mailbox coverage from NIP-65" ); self.retire_idle_nip65_discovery_source(&result.source_relay) @@ -5532,26 +6082,36 @@ impl SyncManager { return; } let now = Instant::now(); - let has_due_author = self + let has_author_work = + self.nip65_discovery + .author_sources + .iter() + .any(|(author, sources)| { + sources.contains(source) + && (self + .nip65_discovery + .in_flight + .contains(&(source.to_string(), *author)) + || self + .nip65_discovery + .next_attempt_at + .get(&(source.to_string(), *author)) + .is_none_or(|due| *due <= now)) + }); + let has_mailbox_work = self .nip65_discovery - .author_sources - .iter() - .any(|(author, sources)| { - sources.contains(source) - && !self - .nip65_discovery - .in_flight - .contains(&(source.to_string(), *author)) - && self - .nip65_discovery - .next_attempt_at - .get(&(source.to_string(), *author)) - .is_none_or(|due| *due <= now) - }); - if !has_due_author { + .mailbox_probes_in_flight + .contains(source) + || (self.nip65_discovery.mailbox_roots.contains_key(source) + && self + .nip65_discovery + .mailbox_probe_next_at + .get(source) + .is_none_or(|due| *due <= now)); + if !has_author_work && !has_mailbox_work && !self.has_pending_batches(source).await { tracing::debug!( relay = %source, - "Retiring idle NIP-65 discovery connection" + "Retiring idle NIP-65 control-plane connection" ); self.disconnect_relay(source).await; } @@ -5627,6 +6187,15 @@ impl SyncManager { if let Some(ref metrics) = self.metrics { metrics.record_connection_attempt(&result.relay_url, true); } + if self + .nip65_discovery + .mailbox_roots + .contains_key(&result.relay_url) + { + self.nip65_discovery + .mailbox_probe_next_at + .insert(result.relay_url.clone(), Instant::now()); + } self.handle_connect_or_reconnect(&result.relay_url).await; } ConnectAttemptOutcome::Failed(error) => { @@ -8492,6 +9061,8 @@ mod tests { #[tokio::test] async fn descendant_frontier_recurses_through_event_and_coordinate_references() { let keys = Keys::generate(); + let participant = Keys::generate(); + let reactor = Keys::generate(); let root = EventBuilder::new(Kind::GitIssue, "root") .finalize(&keys) .expect("build root"); @@ -8502,11 +9073,11 @@ mod tests { let coordinate = descendant_event_coordinate(&addressable).unwrap(); let coordinate_child = EventBuilder::new(Kind::TextNote, "coordinate child") .tag(Tag::custom("a", [coordinate.clone()])) - .finalize(&keys) + .finalize(&participant) .expect("build coordinate child"); let grandchild = EventBuilder::new(Kind::TextNote, "grandchild") .tag(Tag::event(coordinate_child.id)) - .finalize(&keys) + .finalize(&reactor) .expect("build grandchild"); let unrelated = EventBuilder::new(Kind::TextNote, "unrelated") .finalize(&keys) @@ -8530,9 +9101,73 @@ mod tests { HashSet::from([addressable.id, coordinate_child.id, grandchild.id]) ); assert_eq!(frontier.coordinates, HashSet::from([coordinate])); + assert_eq!( + frontier.author_roots, + HashMap::from([ + (keys.public_key(), HashSet::from([root.id])), + (participant.public_key(), HashSet::from([root.id])), + (reactor.public_key(), HashSet::from([root.id])), + ]) + ); assert!(!frontier.event_ids.contains(&unrelated.id)); } + #[tokio::test] + async fn descendant_frontier_preserves_exact_root_provenance() { + let owner = Keys::generate(); + let first_participant = Keys::generate(); + let second_participant = Keys::generate(); + let shared_participant = Keys::generate(); + let first_root = EventBuilder::new(Kind::GitIssue, "first root") + .finalize(&owner) + .expect("build first root"); + let second_root = EventBuilder::new(Kind::GitIssue, "second root") + .finalize(&owner) + .expect("build second root"); + let first_reply = EventBuilder::new(Kind::TextNote, "first reply") + .tag(Tag::event(first_root.id)) + .finalize(&first_participant) + .expect("build first reply"); + let second_reply = EventBuilder::new(Kind::TextNote, "second reply") + .tag(Tag::event(second_root.id)) + .finalize(&second_participant) + .expect("build second reply"); + let shared_reply = EventBuilder::new(Kind::TextNote, "shared reply") + .tags([Tag::event(first_reply.id), Tag::event(second_reply.id)]) + .finalize(&shared_participant) + .expect("build shared reply"); + let database: SharedDatabase = Arc::new(nostr_memory::MemoryDatabase::unbounded()); + for event in [ + &first_root, + &second_root, + &first_reply, + &second_reply, + &shared_reply, + ] { + database.save_event(event).await.expect("save test event"); + } + + let frontier = recursive_descendant_frontier( + &database, + &HashSet::from([first_root.id, second_root.id]), + 500, + ) + .await; + + assert_eq!( + frontier.author_roots[&first_participant.public_key()], + HashSet::from([first_root.id]) + ); + assert_eq!( + frontier.author_roots[&second_participant.public_key()], + HashSet::from([second_root.id]) + ); + assert_eq!( + frontier.author_roots[&shared_participant.public_key()], + HashSet::from([first_root.id, second_root.id]) + ); + } + #[tokio::test] async fn descendant_frontier_stops_at_the_depth_bound() { let keys = Keys::generate(); @@ -8636,6 +9271,7 @@ mod tests { let frontier = DescendantFrontier { event_ids: HashSet::from([EventId::from_byte_array([7; 32])]), coordinates: HashSet::from([format!("30023:{}:article", "a".repeat(64))]), + ..DescendantFrontier::default() }; let filters = descendant_frontier_filters(&frontier, Some(since)); @@ -8660,6 +9296,138 @@ mod tests { ); } + #[test] + fn mailbox_probe_filters_cover_repository_roots_and_descendants() { + let root = EventId::from_byte_array([21; 32]); + let descendant = EventId::from_byte_array([22; 32]); + let repository = "30617:owner:repo".to_string(); + let coordinate = "30023:participant:thread".to_string(); + let frontier = DescendantFrontier { + event_ids: HashSet::from([descendant]), + coordinates: HashSet::from([coordinate.clone()]), + ..Default::default() + }; + + let filters = mailbox_probe_filters( + &HashSet::from([repository.clone()]), + &HashSet::from([root]), + &frontier, + ); + let serialized = filters.iter().map(Filter::as_json).collect::>(); + + assert_eq!(filters.len(), 6, "a/A/q and e/E/q must be probed"); + for value in [repository, coordinate, root.to_hex(), descendant.to_hex()] { + assert!(serialized.iter().any(|filter| filter.contains(&value))); + } + } + + #[test] + fn mailbox_probe_order_prioritizes_longest_waiting_due_source() { + let now = Instant::now(); + let root = EventId::from_byte_array([23; 32]); + let roots = HashMap::from([ + ("wss://old.example".to_string(), HashSet::from([root])), + ( + "wss://large.example".to_string(), + HashSet::from([ + root, + EventId::from_byte_array([24; 32]), + EventId::from_byte_array([25; 32]), + ]), + ), + ("wss://future.example".to_string(), HashSet::from([root])), + ]); + let next_at = HashMap::from([ + ( + "wss://old.example".to_string(), + now - Duration::from_secs(10), + ), + ( + "wss://large.example".to_string(), + now - Duration::from_secs(1), + ), + ( + "wss://future.example".to_string(), + now + Duration::from_secs(1), + ), + ]); + + assert_eq!( + due_mailbox_relays(&roots, &next_at, now), + vec![ + "wss://old.example".to_string(), + "wss://large.example".to_string(), + ] + ); + } + + #[test] + fn mailbox_probe_waits_for_committed_connection_lifecycle() { + assert!(!mailbox_probe_connection_ready( + Some(ConnectionStatus::Connecting), + true, + )); + assert!(mailbox_probe_connection_ready( + Some(ConnectionStatus::Syncing), + true, + )); + assert!(!mailbox_probe_connection_ready( + Some(ConnectionStatus::Connected), + false, + )); + } + + #[test] + fn mailbox_probe_prefers_ready_relay_over_older_unavailable_relay() { + let unavailable = "wss://unavailable.example".to_string(); + let ready = "wss://ready.example".to_string(); + let due = vec![unavailable.clone(), ready.clone()]; + + assert_eq!( + select_due_mailbox_relay(&due, &HashSet::from([ready.clone()])), + Some(ready), + ); + assert_eq!( + select_due_mailbox_relay(&due, &HashSet::new()), + Some(unavailable), + ); + } + + #[test] + fn mailbox_cursor_survives_inventory_growth_and_retires_with_source() { + let relay = "wss://mailbox.example".to_string(); + let first_root = EventId::from_byte_array([31; 32]); + let second_root = EventId::from_byte_array([32; 32]); + let mut discovery = Nip65DiscoveryState::default(); + discovery + .mailbox_roots + .insert(relay.clone(), HashSet::from([first_root])); + discovery.mailbox_probe_next_filter.insert(relay.clone(), 3); + + discovery.install_mailbox_overlay( + HashMap::from([(relay.clone(), HashSet::from([first_root, second_root]))]), + Instant::now(), + ); + assert_eq!(discovery.mailbox_probe_next_filter[&relay], 3); + + discovery.install_mailbox_overlay(HashMap::new(), Instant::now()); + assert!(!discovery.mailbox_probe_next_filter.contains_key(&relay)); + assert!(!discovery.mailbox_probe_next_at.contains_key(&relay)); + } + + #[test] + fn mailbox_completion_advances_after_failure_without_ending_cycle() { + let (next_filter, completed_cycle, retry_after) = mailbox_probe_completion(0, 2, false); + assert_eq!(next_filter, 1); + assert!(!completed_cycle); + assert_eq!(retry_after, mailbox_probe_retry_interval()); + + let (next_filter, completed_cycle, refresh_after) = mailbox_probe_completion(1, 2, true); + assert_eq!(next_filter, 0); + assert!(completed_cycle); + assert_eq!(refresh_after, mailbox_probe_refresh_interval()); + } + #[test] fn own_relay_targets_are_excluded_from_sync_actions() { assert!(is_own_sync_target("wss://gitnostr.com", "gitnostr.com")); @@ -8694,6 +9462,7 @@ mod tests { let frontier = DescendantFrontier { event_ids: HashSet::from([member]), coordinates: HashSet::new(), + ..DescendantFrontier::default() }; let now = Timestamp::from_secs(200_000); let mut rotation = DescendantSyncRotation::default(); @@ -8786,6 +9555,7 @@ mod tests { let frontier = DescendantFrontier { event_ids: HashSet::from([descendant]), coordinates: HashSet::from(["1621:author:patch".to_string()]), + ..DescendantFrontier::default() }; let packed = packed_historic_filters(&items, &frontier); @@ -8816,6 +9586,7 @@ mod tests { &DescendantFrontier { event_ids: members, coordinates: HashSet::new(), + ..DescendantFrontier::default() }, None, ); diff --git a/src/sync/self_subscriber.rs b/src/sync/self_subscriber.rs index b9a5327..9088b63 100644 --- a/src/sync/self_subscriber.rs +++ b/src/sync/self_subscriber.rs @@ -670,11 +670,9 @@ impl SelfSubscriber { // state, avoiding an O(all repositories × all relays) rebuild for // every historic batch. for relay_url in dirty_relays { - // Skip our own relay URL (we're subscribed to ourselves via self-subscription) - if relay_url.contains(&self.relay_domain) { - continue; - } - + // Keep our own relay's dirty signal: SyncManager rejects it as a + // sync target, but uses the signal to refresh participant mailbox + // inventory after a newly accepted local root arrives. let action = AddFilters { relay_url: relay_url.clone(), items: crate::sync::PendingItems::default(), diff --git a/tests/sync/proactive_sync_plus.rs b/tests/sync/proactive_sync_plus.rs index 5c9d2c2..a5c881f 100644 --- a/tests/sync/proactive_sync_plus.rs +++ b/tests/sync/proactive_sync_plus.rs @@ -1,5 +1,6 @@ //! Minimal GRASP-03 mailbox discovery scenario. +use std::path::Path; use std::time::Duration; use nostr_sdk::prelude::*; @@ -9,6 +10,23 @@ use crate::common::{ setup_announcement_on_relay, wait_for_event_on_relay, MockRelay, TestClient, TestRelay, }; +async fn wait_for_log_line(path: &Path, timeout: Duration, predicate: F) -> bool +where + F: Fn(&str) -> bool, +{ + let deadline = tokio::time::Instant::now() + timeout; + loop { + let log = std::fs::read_to_string(path).unwrap_or_default(); + if log.lines().any(&predicate) { + return true; + } + if tokio::time::Instant::now() >= deadline { + return false; + } + tokio::time::sleep(Duration::from_millis(100)).await; + } +} + #[tokio::test] async fn root_author_inbox_reuses_existing_root_sync_pipeline() { let index = MockRelay::start().await; @@ -207,6 +225,91 @@ async fn root_author_inbox_reuses_existing_root_sync_pipeline() { inbox.stop().await; } +#[tokio::test] +async fn participant_write_mailbox_fetches_child_of_direct_reaction() { + let index = MockRelay::start().await; + let mailbox = MockRelay::start().await; + let owner = Keys::generate(); + let root_author = Keys::generate(); + let participant = Keys::generate(); + let child_author = Keys::generate(); + let identifier = "participant-write-mailbox"; + + let relay_list = EventBuilder::new(Kind::RelayList, "") + .tag(Tag::custom("r", vec![mailbox.url(), "write"])) + .finalize(&participant) + .expect("build participant relay list"); + send_to_relay_url(index.url(), &relay_list) + .await + .expect("seed participant relay list"); + + let syncing_git_dir = tempfile::tempdir().expect("create persistent git directory"); + let syncing_relay_dir = tempfile::tempdir().expect("create persistent relay directory"); + let syncing = TestRelay::start_on_reservation_persistent_sync( + reserve_port(), + Some(index.url().to_string()), + false, + syncing_git_dir.path().to_path_buf(), + syncing_relay_dir.path().to_path_buf(), + ) + .await; + let syncing_domain = syncing.domain(); + let (_announcement, _git_dir) = + setup_announcement_on_relay(&syncing, &owner, &[&syncing_domain], identifier).await; + let issue = build_layer2_issue_event( + &root_author, + &repo_coord(&owner, identifier), + "root with a reaction whose child is mailbox-only", + ) + .expect("build accepted root"); + let root_client = TestClient::new(syncing.url(), root_author) + .await + .expect("connect root author"); + root_client.send_event(&issue).await.expect("publish root"); + + let reaction = EventBuilder::new(Kind::Reaction, "+") + .tag(Tag::custom("e", vec![issue.id.to_hex()])) + .finalize(&participant) + .expect("build direct participant reaction"); + let participant_client = TestClient::new(syncing.url(), participant) + .await + .expect("connect participant"); + participant_client + .send_event(&reaction) + .await + .expect("publish direct reaction"); + + let child = EventBuilder::new(Kind::TextNote, "reply to a reaction") + .tag(Tag::custom("e", vec![reaction.id.to_hex()])) + .finalize(&child_author) + .expect("build mailbox-only child"); + send_to_relay_url(mailbox.url(), &child) + .await + .expect("seed participant write mailbox"); + + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(child.id), + Duration::from_secs(30), + ) + .await, + "history probing should fetch a child that only names the participant's direct reaction" + ); + let sync_log = std::fs::read_to_string(syncing.log_path()).expect("read syncing relay log"); + assert!( + sync_log.lines().any(|line| { + line.contains("Started bounded participant mailbox fetch") + && line.contains(mailbox.url()) + }), + "participant mailbox transport should remain observable" + ); + + syncing.stop().await; + mailbox.stop().await; + index.stop().await; +} + #[tokio::test] async fn missing_relay_list_uses_bounded_fallback_coverage() { let index = MockRelay::start().await; @@ -259,7 +362,7 @@ async fn missing_relay_list_uses_bounded_fallback_coverage() { let sync_log = std::fs::read_to_string(syncing.log_path()).expect("read syncing relay log"); assert!( sync_log.lines().any(|line| { - line.contains("Updated proactive inbox coverage from NIP-65") + line.contains("Updated proactive participant mailbox coverage from NIP-65") && line.contains("fallback_authors=1") }), "fallback activation should remain observable" @@ -288,12 +391,11 @@ async fn missing_relay_list_uses_bounded_fallback_coverage() { .await, "a later accepted relay list should replace desired fallback coverage" ); - let sync_log = std::fs::read_to_string(syncing.log_path()).expect("read updated relay log"); assert!( - sync_log.lines().any(|line| { - line.contains("Updated proactive inbox coverage from NIP-65") - && line.contains("fallback_authors=0") - }), + wait_for_log_line(&syncing.log_path(), Duration::from_secs(20), |line| { + line.contains("proactive participant mailbox") && line.contains("fallback_authors=0") + }) + .await, "accepted NIP-65 ownership should retire the author from desired fallback coverage" );