mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
fix(sync): backfill repository mailbox history
Historical repository roots and root-only descendants can exist on an accepted owner or maintainer's NIP-65 mailbox without appearing on a declared repository relay. Root-participant provenance alone cannot discover that first root or query a known root against an unrelated maintainer mailbox. Associate accepted Full repositories with their owners and declared maintainers, derive a repository-scoped history overlay from retained relay lists, and add every known root in that repository to the mailbox rotation. Reuse the existing byte-bounded, paced, single-flight fetch path and reschedule when root inventory grows. Correctness depends on repository admission remaining authoritative and on every fetched event continuing through the normal write policy, persistence pipeline, and global NIP-09/NIP-62 tombstones. Participant-only mailboxes retain exact root provenance. Deliberately excluded are permanent maintainer live subscriptions, repository-coordinate expansion in private mode, and any change to deletion or retention behavior. Validated with nix develop -c cargo test, nix develop -c cargo fmt -- --check, and nix develop -c cargo clippy --all-targets --all-features -- -D warnings.
This commit is contained in:
@@ -137,6 +137,13 @@ Performance and Security fixes - along with other improvements; immediate upgrad
|
||||
|
||||
### Fixed
|
||||
|
||||
- Discover historical repository roots and descendants from the NIP-65
|
||||
mailboxes of accepted repository owners and maintainers. Public Sync+
|
||||
instances query exact repository coordinates and every known root for that
|
||||
repository through the existing paced, byte-bounded history workers, while
|
||||
private instances continue to withhold repository coordinates and the
|
||||
ordinary write policy and persistent deletion tombstones remain authoritative.
|
||||
|
||||
- Retire peer-closed outbound subscriptions from rust-nostr's desired registry
|
||||
without interrupting NIP-42: the first `auth-required` response keeps the
|
||||
same subscription and ledger slot for one authenticated retry, while a repeat
|
||||
|
||||
@@ -685,7 +685,13 @@ 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.
|
||||
subscriptions or coupling progress between relays. On public instances, those
|
||||
workers also query each accepted repository owner's and declared maintainer's
|
||||
mailboxes for their exact Full repository coordinates and all locally known
|
||||
roots in those repositories. This discovers previously unknown roots and
|
||||
descendants that exist only on a maintainer mailbox without broadening an
|
||||
ordinary participant mailbox. Private instances omit repository-coordinate
|
||||
mailbox expansion.
|
||||
|
||||
### Rejected Events Index
|
||||
|
||||
|
||||
@@ -18,7 +18,10 @@ 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.
|
||||
repository and thread-reference filters. On public instances, the same
|
||||
history-only probe associates accepted repository owners and declared
|
||||
maintainers with those repositories and queries their mailboxes by repository
|
||||
coordinate. This can discover a root which is not yet present locally.
|
||||
|
||||
Mailbox work reuses GRASP-02's connection safety, filter byte packing,
|
||||
pagination, subscription ledger, request pacing, event pipeline and write
|
||||
@@ -96,14 +99,30 @@ 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
|
||||
`e`/`E`/`q`. A repository-scoped owner or maintainer mailbox also receives
|
||||
event-reference filters for every currently known root in that repository.
|
||||
This permits a status or other descendant stored only on that mailbox to be
|
||||
found even when the root author's current relay list does not name it. A
|
||||
participant-only mailbox remains constrained to its derived roots. A numeric
|
||||
in-memory cursor gives each byte-bounded 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.
|
||||
|
||||
Owner/maintainer expansion can increase a relay's history rotation in
|
||||
proportion to the byte-packed repository coordinates and known roots assigned
|
||||
to it. Shared relays and duplicate repository scopes are unioned before filter
|
||||
construction. The expansion does not increase instantaneous worker
|
||||
concurrency on any one relay: the manager still starts at most one due relay
|
||||
per maintenance pass, each relay remains single-flight, and each worker fetches
|
||||
one filter group through the existing background request pacer before
|
||||
yielding. More distinct maintainer relays can nevertheless produce more active
|
||||
per-relay workers across successive passes; the fixed-cardinality relay,
|
||||
cursor, and active-worker metrics expose that growth for a production soak.
|
||||
|
||||
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
|
||||
@@ -127,8 +146,11 @@ events on restart; this change does not add a second deletion graph.
|
||||
## Deliberately excluded
|
||||
|
||||
- mailbox expansion for authors who have not produced locally accepted
|
||||
repository-thread events;
|
||||
- maintainer mailbox expansion unrelated to an accepted root;
|
||||
repository-thread events and are not an accepted repository owner or
|
||||
maintainer;
|
||||
- repository mailbox expansion beyond an owner's or declared maintainer's
|
||||
exact accepted Full repositories;
|
||||
- repository-coordinate mailbox expansion in private mode;
|
||||
- identity storage for authors outside accepted repositories and their
|
||||
locally accepted threads;
|
||||
- per-user fallback configuration knobs; and
|
||||
|
||||
+131
-9
@@ -140,7 +140,20 @@ pub fn accepted_repository_authors<'a>(
|
||||
announcements: impl IntoIterator<Item = &'a Event>,
|
||||
repositories: &HashMap<String, RepoSyncNeeds>,
|
||||
) -> HashSet<PublicKey> {
|
||||
let mut authors = HashSet::new();
|
||||
accepted_repository_author_repositories(announcements, repositories)
|
||||
.into_keys()
|
||||
.collect()
|
||||
}
|
||||
|
||||
/// Associate owners and declared maintainers with each accepted Full
|
||||
/// repository. The caller may use this for repository-scoped mailbox history;
|
||||
/// keeping the coordinate provenance here prevents one maintainer's relay list
|
||||
/// from broadening coverage to unrelated repositories.
|
||||
pub fn accepted_repository_author_repositories<'a>(
|
||||
announcements: impl IntoIterator<Item = &'a Event>,
|
||||
repositories: &HashMap<String, RepoSyncNeeds>,
|
||||
) -> HashMap<PublicKey, HashSet<String>> {
|
||||
let mut author_repositories: HashMap<PublicKey, HashSet<String>> = HashMap::new();
|
||||
for event in announcements {
|
||||
if event.kind != Kind::GitRepoAnnouncement {
|
||||
continue;
|
||||
@@ -159,7 +172,10 @@ pub fn accepted_repository_authors<'a>(
|
||||
{
|
||||
continue;
|
||||
}
|
||||
authors.insert(event.pubkey);
|
||||
author_repositories
|
||||
.entry(event.pubkey)
|
||||
.or_default()
|
||||
.insert(repository.clone());
|
||||
// Include every currently-listed maintainer (active `M`/`m` role tags
|
||||
// or the deprecated maintainers fallback). Invited pubkeys are
|
||||
// included deliberately: their announcements must be fetched so we
|
||||
@@ -167,15 +183,19 @@ pub fn accepted_repository_authors<'a>(
|
||||
if let Ok(announcement) =
|
||||
crate::nostr::events::RepositoryAnnouncement::from_event(event.clone())
|
||||
{
|
||||
authors.extend(
|
||||
announcement
|
||||
.listed_maintainers()
|
||||
.iter()
|
||||
.filter_map(|value| PublicKey::from_hex(value).ok()),
|
||||
);
|
||||
for author in announcement
|
||||
.listed_maintainers()
|
||||
.iter()
|
||||
.filter_map(|value| PublicKey::from_hex(value).ok())
|
||||
{
|
||||
author_repositories
|
||||
.entry(author)
|
||||
.or_default()
|
||||
.insert(repository.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
authors
|
||||
author_repositories
|
||||
}
|
||||
|
||||
/// Associate each accepted root author with their own roots and each accepted
|
||||
@@ -242,6 +262,42 @@ pub fn build_inbox_root_overlay(
|
||||
overlay
|
||||
}
|
||||
|
||||
/// Build the history-only repository scope advertised by accepted repository
|
||||
/// owners and maintainers. This is deliberately separate from root provenance:
|
||||
/// repository scope may discover roots not yet present locally, while a thread
|
||||
/// participant remains constrained to roots in which they actually appeared.
|
||||
pub fn build_mailbox_repository_overlay(
|
||||
author_repositories: &HashMap<PublicKey, HashSet<String>>,
|
||||
author_mailboxes: &HashMap<PublicKey, HashSet<String>>,
|
||||
fallback_authors: &HashSet<PublicKey>,
|
||||
fallback_relays: &HashSet<String>,
|
||||
) -> HashMap<String, HashSet<String>> {
|
||||
let mut overlay: HashMap<String, HashSet<String>> = HashMap::new();
|
||||
for (author, mailboxes) in author_mailboxes {
|
||||
let Some(repositories) = author_repositories.get(author) else {
|
||||
continue;
|
||||
};
|
||||
for mailbox in mailboxes {
|
||||
overlay
|
||||
.entry(mailbox.clone())
|
||||
.or_default()
|
||||
.extend(repositories.iter().cloned());
|
||||
}
|
||||
}
|
||||
for author in fallback_authors {
|
||||
let Some(repositories) = author_repositories.get(author) else {
|
||||
continue;
|
||||
};
|
||||
for relay in fallback_relays {
|
||||
overlay
|
||||
.entry(relay.clone())
|
||||
.or_default()
|
||||
.extend(repositories.iter().cloned());
|
||||
}
|
||||
}
|
||||
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<String> {
|
||||
@@ -408,6 +464,22 @@ mod tests {
|
||||
accepted_repository_authors([&accepted, &state_only], &repositories),
|
||||
HashSet::from([owner.public_key(), maintainer.public_key()])
|
||||
);
|
||||
let author_repositories =
|
||||
accepted_repository_author_repositories([&accepted, &state_only], &repositories);
|
||||
let accepted_repository = format!("30617:{}:{identifier}", owner.public_key());
|
||||
assert_eq!(
|
||||
author_repositories,
|
||||
HashMap::from([
|
||||
(
|
||||
owner.public_key(),
|
||||
HashSet::from([accepted_repository.clone()]),
|
||||
),
|
||||
(
|
||||
maintainer.public_key(),
|
||||
HashSet::from([accepted_repository]),
|
||||
),
|
||||
])
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -575,6 +647,56 @@ mod tests {
|
||||
.all(|roots| roots == &HashSet::from([root])));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn maintainer_mailboxes_receive_only_their_accepted_repositories() {
|
||||
let maintainer = Keys::generate().public_key();
|
||||
let unrelated = Keys::generate().public_key();
|
||||
let repository = "30617:owner:accepted".to_string();
|
||||
let overlay = build_mailbox_repository_overlay(
|
||||
&HashMap::from([(maintainer, HashSet::from([repository.clone()]))]),
|
||||
&HashMap::from([
|
||||
(
|
||||
maintainer,
|
||||
HashSet::from(["wss://maintainer.example".to_string()]),
|
||||
),
|
||||
(
|
||||
unrelated,
|
||||
HashSet::from(["wss://unrelated.example".to_string()]),
|
||||
),
|
||||
]),
|
||||
&HashSet::new(),
|
||||
&HashSet::new(),
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
overlay,
|
||||
HashMap::from([(
|
||||
"wss://maintainer.example".to_string(),
|
||||
HashSet::from([repository]),
|
||||
)])
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_maintainer_lists_use_the_bounded_fallback_scope() {
|
||||
let maintainer = Keys::generate().public_key();
|
||||
let repository = "30617:owner:accepted".to_string();
|
||||
let overlay = build_mailbox_repository_overlay(
|
||||
&HashMap::from([(maintainer, HashSet::from([repository.clone()]))]),
|
||||
&HashMap::new(),
|
||||
&HashSet::from([maintainer]),
|
||||
&HashSet::from(["wss://fallback.example".to_string()]),
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
overlay,
|
||||
HashMap::from([(
|
||||
"wss://fallback.example".to_string(),
|
||||
HashSet::from([repository]),
|
||||
)])
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn identity_selection_keeps_latest_profile_and_relay_list_per_author() {
|
||||
let author = Keys::generate();
|
||||
|
||||
+263
-43
@@ -87,21 +87,28 @@ fn mailbox_probe_retry_interval() -> Duration {
|
||||
|
||||
fn due_mailbox_relays(
|
||||
mailbox_roots: &HashMap<String, HashSet<EventId>>,
|
||||
mailbox_repositories: &HashMap<String, HashSet<String>>,
|
||||
next_at: &HashMap<String, Instant>,
|
||||
now: Instant,
|
||||
) -> Vec<String> {
|
||||
let mut relays: Vec<String> = mailbox_roots
|
||||
let mut relays: HashSet<String> = mailbox_roots
|
||||
.keys()
|
||||
.filter(|relay| next_at.get(*relay).is_none_or(|due| *due <= now))
|
||||
.chain(mailbox_repositories.keys())
|
||||
.cloned()
|
||||
.collect();
|
||||
relays.retain(|relay| next_at.get(relay).is_none_or(|due| *due <= now));
|
||||
let mut relays: Vec<String> = relays.into_iter().collect();
|
||||
relays.sort_by(|left, right| {
|
||||
let workload = |relay: &String| {
|
||||
mailbox_roots.get(relay).map_or(0, HashSet::len)
|
||||
+ mailbox_repositories.get(relay).map_or(0, HashSet::len)
|
||||
};
|
||||
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(|| workload(right).cmp(&workload(left)))
|
||||
.then_with(|| left.cmp(right))
|
||||
});
|
||||
relays
|
||||
@@ -125,6 +132,17 @@ fn select_due_mailbox_relay(
|
||||
.cloned()
|
||||
}
|
||||
|
||||
fn public_repository_mailbox_scope(
|
||||
private_mode: bool,
|
||||
author_repositories: HashMap<PublicKey, HashSet<String>>,
|
||||
) -> HashMap<PublicKey, HashSet<String>> {
|
||||
if private_mode {
|
||||
HashMap::new()
|
||||
} else {
|
||||
author_repositories
|
||||
}
|
||||
}
|
||||
|
||||
fn mailbox_probe_completion(
|
||||
filter_index: usize,
|
||||
filter_count: usize,
|
||||
@@ -1481,6 +1499,46 @@ fn mailbox_probe_filters(
|
||||
filters
|
||||
}
|
||||
|
||||
fn expand_mailbox_probe_scope(
|
||||
repository_scope: &HashSet<String>,
|
||||
root_scope: &HashSet<EventId>,
|
||||
root_repositories: &HashMap<EventId, String>,
|
||||
) -> (HashSet<String>, HashSet<EventId>) {
|
||||
let mut repositories = repository_scope.clone();
|
||||
repositories.extend(
|
||||
root_scope
|
||||
.iter()
|
||||
.filter_map(|root| root_repositories.get(root).cloned()),
|
||||
);
|
||||
|
||||
let mut roots = root_scope.clone();
|
||||
roots.extend(
|
||||
root_repositories.iter().filter_map(|(root, repository)| {
|
||||
repository_scope.contains(repository).then_some(*root)
|
||||
}),
|
||||
);
|
||||
(repositories, roots)
|
||||
}
|
||||
|
||||
fn merge_repository_roots_into_mailbox_overlay(
|
||||
mailbox_roots: &mut HashMap<String, HashSet<EventId>>,
|
||||
mailbox_repositories: &HashMap<String, HashSet<String>>,
|
||||
root_repositories: &HashMap<EventId, String>,
|
||||
) {
|
||||
for (relay, repositories) in mailbox_repositories {
|
||||
let roots: HashSet<EventId> = root_repositories
|
||||
.iter()
|
||||
.filter_map(|(root, repository)| repositories.contains(repository).then_some(*root))
|
||||
.collect();
|
||||
if !roots.is_empty() {
|
||||
mailbox_roots
|
||||
.entry(relay.clone())
|
||||
.or_default()
|
||||
.extend(roots);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn packed_historic_filters(items: &PendingItems, frontier: &DescendantFrontier) -> Vec<Filter> {
|
||||
let all_repos: HashSet<_> = items
|
||||
.repos
|
||||
@@ -1860,6 +1918,10 @@ struct Nip65DiscoveryState {
|
||||
/// the paced history-only mailbox path below.
|
||||
root_author_roots: HashMap<PublicKey, HashSet<EventId>>,
|
||||
author_roots: HashMap<PublicKey, HashSet<EventId>>,
|
||||
/// Public accepted repositories owned or maintained by each author. This
|
||||
/// remains empty in private mode so mailbox probes never disclose private
|
||||
/// repository coordinates to a public NIP-65 relay.
|
||||
author_repositories: HashMap<PublicKey, HashSet<String>>,
|
||||
root_repositories: HashMap<EventId, String>,
|
||||
eligible_authors: HashSet<PublicKey>,
|
||||
author_sources: HashMap<PublicKey, HashSet<String>>,
|
||||
@@ -1872,6 +1934,7 @@ struct Nip65DiscoveryState {
|
||||
fallback_authors: HashSet<PublicKey>,
|
||||
inbox_roots: HashMap<String, HashSet<EventId>>,
|
||||
mailbox_roots: HashMap<String, HashSet<EventId>>,
|
||||
mailbox_repositories: HashMap<String, HashSet<String>>,
|
||||
mailbox_probe_next_at: HashMap<String, Instant>,
|
||||
mailbox_probe_next_filter: HashMap<String, usize>,
|
||||
/// At most one ordinary fetch runs per mailbox relay. Different relays
|
||||
@@ -1883,25 +1946,49 @@ struct Nip65DiscoveryState {
|
||||
}
|
||||
|
||||
impl Nip65DiscoveryState {
|
||||
fn has_mailbox_scope(&self, relay: &str) -> bool {
|
||||
self.mailbox_roots.contains_key(relay) || self.mailbox_repositories.contains_key(relay)
|
||||
}
|
||||
|
||||
fn mailbox_relay_count(&self) -> usize {
|
||||
self.mailbox_roots.len()
|
||||
+ self
|
||||
.mailbox_repositories
|
||||
.keys()
|
||||
.filter(|relay| !self.mailbox_roots.contains_key(*relay))
|
||||
.count()
|
||||
}
|
||||
|
||||
fn install_mailbox_overlay(
|
||||
&mut self,
|
||||
new_overlay: HashMap<String, HashSet<EventId>>,
|
||||
new_root_overlay: HashMap<String, HashSet<EventId>>,
|
||||
new_repository_overlay: HashMap<String, HashSet<String>>,
|
||||
now: Instant,
|
||||
) {
|
||||
let old_overlay = std::mem::replace(&mut self.mailbox_roots, new_overlay);
|
||||
if old_overlay == self.mailbox_roots {
|
||||
let old_root_overlay = std::mem::replace(&mut self.mailbox_roots, new_root_overlay);
|
||||
let old_repository_overlay =
|
||||
std::mem::replace(&mut self.mailbox_repositories, new_repository_overlay);
|
||||
if old_root_overlay == self.mailbox_roots
|
||||
&& old_repository_overlay == self.mailbox_repositories
|
||||
{
|
||||
return;
|
||||
}
|
||||
let current_relays: HashSet<String> = self
|
||||
.mailbox_roots
|
||||
.keys()
|
||||
.chain(self.mailbox_repositories.keys())
|
||||
.cloned()
|
||||
.collect();
|
||||
self.mailbox_probe_next_at
|
||||
.retain(|relay, _| self.mailbox_roots.contains_key(relay));
|
||||
.retain(|relay, _| current_relays.contains(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) {
|
||||
.retain(|relay, _| current_relays.contains(relay));
|
||||
for relay in current_relays {
|
||||
if old_root_overlay.get(&relay) != self.mailbox_roots.get(&relay)
|
||||
|| old_repository_overlay.get(&relay) != self.mailbox_repositories.get(&relay)
|
||||
{
|
||||
self.mailbox_probe_next_at.insert(relay.clone(), now);
|
||||
self.mailbox_probe_next_filter
|
||||
.entry(relay.clone())
|
||||
.or_default();
|
||||
self.mailbox_probe_next_filter.entry(relay).or_default();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1916,7 +2003,7 @@ impl Nip65DiscoveryState {
|
||||
// A relay removed from the overlay while its probe was in flight must
|
||||
// not leave orphan cursor or due-time entries; re-addition is reseeded
|
||||
// by `install_mailbox_overlay`.
|
||||
if !self.mailbox_roots.contains_key(relay) {
|
||||
if !self.has_mailbox_scope(relay) {
|
||||
return;
|
||||
}
|
||||
self.mailbox_probe_next_filter
|
||||
@@ -2419,7 +2506,7 @@ async fn run_health_and_metrics_checker(
|
||||
),
|
||||
(
|
||||
"mailbox_probe_relays",
|
||||
manager.nip65_discovery.mailbox_roots.len(),
|
||||
manager.nip65_discovery.mailbox_relay_count(),
|
||||
),
|
||||
(
|
||||
"mailbox_probe_cursors",
|
||||
@@ -5626,8 +5713,14 @@ impl SyncManager {
|
||||
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);
|
||||
let author_repositories = discovery::accepted_repository_author_repositories(
|
||||
announcements.iter(),
|
||||
&repo_index,
|
||||
);
|
||||
let mut eligible_authors: HashSet<PublicKey> =
|
||||
author_repositories.keys().copied().collect();
|
||||
let author_repositories =
|
||||
public_repository_mailbox_scope(self.config.private_mode, author_repositories);
|
||||
drop(candidates);
|
||||
drop(repo_index);
|
||||
let root_ids: HashSet<EventId> = roots.iter().map(|root| root.id).collect();
|
||||
@@ -5653,6 +5746,7 @@ impl SyncManager {
|
||||
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.author_repositories = author_repositories;
|
||||
self.nip65_discovery.eligible_authors = eligible_authors;
|
||||
self.proactive_participant_authors
|
||||
.replace(self.nip65_discovery.eligible_authors.clone())
|
||||
@@ -5759,19 +5853,31 @@ impl SyncManager {
|
||||
&self.nip65_discovery.fallback_authors,
|
||||
&fallback_relays,
|
||||
);
|
||||
let mailbox_overlay = discovery::build_inbox_root_overlay(
|
||||
let mut mailbox_overlay = discovery::build_inbox_root_overlay(
|
||||
&self.nip65_discovery.author_roots,
|
||||
&self.nip65_discovery.author_mailboxes,
|
||||
&self.nip65_discovery.fallback_authors,
|
||||
&fallback_relays,
|
||||
);
|
||||
let mailbox_repository_overlay = discovery::build_mailbox_repository_overlay(
|
||||
&self.nip65_discovery.author_repositories,
|
||||
&self.nip65_discovery.author_mailboxes,
|
||||
&self.nip65_discovery.fallback_authors,
|
||||
&fallback_relays,
|
||||
);
|
||||
merge_repository_roots_into_mailbox_overlay(
|
||||
&mut mailbox_overlay,
|
||||
&mailbox_repository_overlay,
|
||||
&self.nip65_discovery.root_repositories,
|
||||
);
|
||||
self.install_nip65_inbox_overlay(inbox_overlay).await;
|
||||
self.install_nip65_mailbox_overlay(mailbox_overlay);
|
||||
self.install_nip65_mailbox_overlay(mailbox_overlay, mailbox_repository_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(),
|
||||
repository_mailbox_author_count = self.nip65_discovery.author_repositories.len(),
|
||||
mailbox_relay_count = self.nip65_discovery.mailbox_relay_count(),
|
||||
fallback_authors = self.nip65_discovery.fallback_authors.len(),
|
||||
"Reconciled proactive participant mailbox inventory"
|
||||
);
|
||||
@@ -5881,9 +5987,16 @@ impl SyncManager {
|
||||
}
|
||||
}
|
||||
|
||||
fn install_nip65_mailbox_overlay(&mut self, new_overlay: HashMap<String, HashSet<EventId>>) {
|
||||
self.nip65_discovery
|
||||
.install_mailbox_overlay(new_overlay, Instant::now());
|
||||
fn install_nip65_mailbox_overlay(
|
||||
&mut self,
|
||||
new_root_overlay: HashMap<String, HashSet<EventId>>,
|
||||
new_repository_overlay: HashMap<String, HashSet<String>>,
|
||||
) {
|
||||
self.nip65_discovery.install_mailbox_overlay(
|
||||
new_root_overlay,
|
||||
new_repository_overlay,
|
||||
Instant::now(),
|
||||
);
|
||||
}
|
||||
|
||||
async fn schedule_mailbox_probe(&mut self) {
|
||||
@@ -5894,6 +6007,7 @@ impl SyncManager {
|
||||
let now = Instant::now();
|
||||
let due_relays: Vec<String> = due_mailbox_relays(
|
||||
&self.nip65_discovery.mailbox_roots,
|
||||
&self.nip65_discovery.mailbox_repositories,
|
||||
&self.nip65_discovery.mailbox_probe_next_at,
|
||||
now,
|
||||
)
|
||||
@@ -5954,11 +6068,23 @@ impl SyncManager {
|
||||
return;
|
||||
}
|
||||
|
||||
let roots = self.nip65_discovery.mailbox_roots[&relay].clone();
|
||||
let repositories: HashSet<String> = roots
|
||||
.iter()
|
||||
.filter_map(|root| self.nip65_discovery.root_repositories.get(root).cloned())
|
||||
.collect();
|
||||
let root_scope = self
|
||||
.nip65_discovery
|
||||
.mailbox_roots
|
||||
.get(&relay)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let repository_scope = self
|
||||
.nip65_discovery
|
||||
.mailbox_repositories
|
||||
.get(&relay)
|
||||
.cloned()
|
||||
.unwrap_or_default();
|
||||
let (repositories, roots) = expand_mailbox_probe_scope(
|
||||
&repository_scope,
|
||||
&root_scope,
|
||||
&self.nip65_discovery.root_repositories,
|
||||
);
|
||||
let frontier = recursive_descendant_frontier(
|
||||
&self.database,
|
||||
&roots,
|
||||
@@ -6010,7 +6136,7 @@ impl SyncManager {
|
||||
}
|
||||
|
||||
fn defer_mailbox_relay(&mut self, relay: &str) {
|
||||
if self.nip65_discovery.mailbox_roots.contains_key(relay) {
|
||||
if self.nip65_discovery.has_mailbox_scope(relay) {
|
||||
self.nip65_discovery.mailbox_probe_next_at.insert(
|
||||
relay.to_string(),
|
||||
Instant::now() + mailbox_probe_retry_interval(),
|
||||
@@ -6134,7 +6260,13 @@ impl SyncManager {
|
||||
.insert((result.source_relay.clone(), *author), now + retry_after);
|
||||
if authors_with_lists.contains(author) {
|
||||
self.nip65_discovery.fallback_authors.remove(author);
|
||||
} else if query_succeeded && self.nip65_discovery.author_roots.contains_key(author) {
|
||||
} else if query_succeeded
|
||||
&& (self.nip65_discovery.author_roots.contains_key(author)
|
||||
|| self
|
||||
.nip65_discovery
|
||||
.author_repositories
|
||||
.contains_key(author))
|
||||
{
|
||||
self.nip65_discovery.fallback_authors.insert(*author);
|
||||
}
|
||||
}
|
||||
@@ -6200,18 +6332,29 @@ impl SyncManager {
|
||||
&self.nip65_discovery.fallback_authors,
|
||||
&fallback_relays,
|
||||
);
|
||||
let mailbox_overlay = discovery::build_inbox_root_overlay(
|
||||
let mut mailbox_overlay = discovery::build_inbox_root_overlay(
|
||||
&self.nip65_discovery.author_roots,
|
||||
&self.nip65_discovery.author_mailboxes,
|
||||
&self.nip65_discovery.fallback_authors,
|
||||
&fallback_relays,
|
||||
);
|
||||
let mailbox_repository_overlay = discovery::build_mailbox_repository_overlay(
|
||||
&self.nip65_discovery.author_repositories,
|
||||
&self.nip65_discovery.author_mailboxes,
|
||||
&self.nip65_discovery.fallback_authors,
|
||||
&fallback_relays,
|
||||
);
|
||||
merge_repository_roots_into_mailbox_overlay(
|
||||
&mut mailbox_overlay,
|
||||
&mailbox_repository_overlay,
|
||||
&self.nip65_discovery.root_repositories,
|
||||
);
|
||||
self.install_nip65_inbox_overlay(inbox_overlay).await;
|
||||
self.install_nip65_mailbox_overlay(mailbox_overlay);
|
||||
self.install_nip65_mailbox_overlay(mailbox_overlay, mailbox_repository_overlay);
|
||||
tracing::info!(
|
||||
source = %result.source_relay,
|
||||
authors = result.authors.len(),
|
||||
mailbox_relays = self.nip65_discovery.mailbox_roots.len(),
|
||||
mailbox_relays = self.nip65_discovery.mailbox_relay_count(),
|
||||
fallback_authors = self.nip65_discovery.fallback_authors.len(),
|
||||
"Updated proactive participant mailbox coverage from NIP-65"
|
||||
);
|
||||
@@ -6250,7 +6393,7 @@ impl SyncManager {
|
||||
.nip65_discovery
|
||||
.mailbox_probes_in_flight
|
||||
.contains(source)
|
||||
|| (self.nip65_discovery.mailbox_roots.contains_key(source)
|
||||
|| (self.nip65_discovery.has_mailbox_scope(source)
|
||||
&& self
|
||||
.nip65_discovery
|
||||
.mailbox_probe_next_at
|
||||
@@ -6367,11 +6510,7 @@ 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)
|
||||
{
|
||||
if self.nip65_discovery.has_mailbox_scope(&result.relay_url) {
|
||||
self.nip65_discovery
|
||||
.mailbox_probe_next_at
|
||||
.insert(result.relay_url.clone(), Instant::now());
|
||||
@@ -9555,6 +9694,60 @@ mod tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repository_mailbox_scope_adds_every_known_root_without_broadening_participants() {
|
||||
let first_root = EventId::from_byte_array([23; 32]);
|
||||
let second_root = EventId::from_byte_array([24; 32]);
|
||||
let participant_root = EventId::from_byte_array([25; 32]);
|
||||
let repository = "30617:owner:maintained".to_string();
|
||||
let participant_repository = "30617:owner:participant".to_string();
|
||||
let root_repositories = HashMap::from([
|
||||
(first_root, repository.clone()),
|
||||
(second_root, repository.clone()),
|
||||
(participant_root, participant_repository.clone()),
|
||||
]);
|
||||
|
||||
let (repositories, roots) = expand_mailbox_probe_scope(
|
||||
&HashSet::from([repository.clone()]),
|
||||
&HashSet::from([participant_root]),
|
||||
&root_repositories,
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
repositories,
|
||||
HashSet::from([repository, participant_repository])
|
||||
);
|
||||
assert_eq!(
|
||||
roots,
|
||||
HashSet::from([first_root, second_root, participant_root])
|
||||
);
|
||||
|
||||
let (_, participant_only_roots) = expand_mailbox_probe_scope(
|
||||
&HashSet::new(),
|
||||
&HashSet::from([participant_root]),
|
||||
&root_repositories,
|
||||
);
|
||||
assert_eq!(participant_only_roots, HashSet::from([participant_root]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repository_mailbox_root_overlay_tracks_inventory_growth() {
|
||||
let relay = "wss://maintainer.example".to_string();
|
||||
let repository = "30617:owner:maintained".to_string();
|
||||
let root = EventId::from_byte_array([26; 32]);
|
||||
let mailbox_repositories =
|
||||
HashMap::from([(relay.clone(), HashSet::from([repository.clone()]))]);
|
||||
let mut mailbox_roots = HashMap::new();
|
||||
|
||||
merge_repository_roots_into_mailbox_overlay(
|
||||
&mut mailbox_roots,
|
||||
&mailbox_repositories,
|
||||
&HashMap::from([(root, repository)]),
|
||||
);
|
||||
|
||||
assert_eq!(mailbox_roots[&relay], HashSet::from([root]));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mailbox_probe_order_prioritizes_longest_waiting_due_source() {
|
||||
let now = Instant::now();
|
||||
@@ -9587,7 +9780,7 @@ mod tests {
|
||||
]);
|
||||
|
||||
assert_eq!(
|
||||
due_mailbox_relays(&roots, &next_at, now),
|
||||
due_mailbox_relays(&roots, &HashMap::new(), &next_at, now),
|
||||
vec![
|
||||
"wss://old.example".to_string(),
|
||||
"wss://large.example".to_string(),
|
||||
@@ -9595,6 +9788,30 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn repository_only_mailbox_scope_is_due() {
|
||||
let now = Instant::now();
|
||||
let relay = "wss://maintainer.example".to_string();
|
||||
let repositories = HashMap::from([(
|
||||
relay.clone(),
|
||||
HashSet::from(["30617:owner:repo".to_string()]),
|
||||
)]);
|
||||
|
||||
assert_eq!(
|
||||
due_mailbox_relays(&HashMap::new(), &repositories, &HashMap::new(), now),
|
||||
vec![relay]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn private_mode_does_not_export_repository_mailbox_scope() {
|
||||
let author = Keys::generate().public_key();
|
||||
let scope = HashMap::from([(author, HashSet::from(["30617:owner:private".to_string()]))]);
|
||||
|
||||
assert_eq!(public_repository_mailbox_scope(false, scope.clone()), scope);
|
||||
assert!(public_repository_mailbox_scope(true, scope).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mailbox_probe_waits_for_committed_connection_lifecycle() {
|
||||
assert!(!mailbox_probe_connection_ready(
|
||||
@@ -9637,14 +9854,17 @@ mod tests {
|
||||
.mailbox_roots
|
||||
.insert(relay.clone(), HashSet::from([first_root]));
|
||||
discovery.mailbox_probe_next_filter.insert(relay.clone(), 3);
|
||||
let now = Instant::now();
|
||||
|
||||
discovery.install_mailbox_overlay(
|
||||
HashMap::from([(relay.clone(), HashSet::from([first_root, second_root]))]),
|
||||
Instant::now(),
|
||||
HashMap::new(),
|
||||
now,
|
||||
);
|
||||
assert_eq!(discovery.mailbox_probe_next_filter[&relay], 3);
|
||||
assert_eq!(discovery.mailbox_probe_next_at[&relay], now);
|
||||
|
||||
discovery.install_mailbox_overlay(HashMap::new(), Instant::now());
|
||||
discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), Instant::now());
|
||||
assert!(!discovery.mailbox_probe_next_filter.contains_key(&relay));
|
||||
assert!(!discovery.mailbox_probe_next_at.contains_key(&relay));
|
||||
}
|
||||
@@ -9666,7 +9886,7 @@ mod tests {
|
||||
now + Duration::from_secs(60)
|
||||
);
|
||||
|
||||
discovery.install_mailbox_overlay(HashMap::new(), now);
|
||||
discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), now);
|
||||
discovery.record_probe_completion(&relay, 2, Duration::from_secs(60), now);
|
||||
assert!(!discovery.mailbox_probe_next_filter.contains_key(&relay));
|
||||
assert!(!discovery.mailbox_probe_next_at.contains_key(&relay));
|
||||
|
||||
@@ -310,6 +310,90 @@ async fn participant_write_mailbox_fetches_child_of_direct_reaction() {
|
||||
index.stop().await;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn owner_mailbox_discovers_unknown_root_and_root_only_status() {
|
||||
let index = MockRelay::start().await;
|
||||
let mailbox = MockRelay::start().await;
|
||||
let owner = Keys::generate();
|
||||
let root_author = Keys::generate();
|
||||
let identifier = "owner-repository-mailbox";
|
||||
|
||||
let relay_list = EventBuilder::new(Kind::RelayList, "")
|
||||
.tag(Tag::custom("r", vec![mailbox.url(), "write"]))
|
||||
.finalize(&owner)
|
||||
.expect("build owner relay list");
|
||||
send_to_relay_url(index.url(), &relay_list)
|
||||
.await
|
||||
.expect("seed owner 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 stored only on the repository owner's mailbox",
|
||||
)
|
||||
.expect("build mailbox-only root");
|
||||
let closed = EventBuilder::new(Kind::from(1632), "closed")
|
||||
.tag(Tag::custom("e", vec![issue.id.to_hex()]))
|
||||
.finalize(&owner)
|
||||
.expect("build root-only closed status");
|
||||
send_to_relay_url(mailbox.url(), &issue)
|
||||
.await
|
||||
.expect("seed unknown root on owner mailbox");
|
||||
send_to_relay_url(mailbox.url(), &closed)
|
||||
.await
|
||||
.expect("seed root-only status on owner mailbox");
|
||||
|
||||
assert!(
|
||||
wait_for_event_on_relay(
|
||||
syncing.url(),
|
||||
Filter::new().id(issue.id),
|
||||
Duration::from_secs(30),
|
||||
)
|
||||
.await,
|
||||
"the owner's repository-coordinate scope should discover an unknown root"
|
||||
);
|
||||
assert!(
|
||||
wait_for_event_on_relay(
|
||||
syncing.url(),
|
||||
Filter::new().id(closed.id),
|
||||
Duration::from_secs(30),
|
||||
)
|
||||
.await,
|
||||
"the owner's mailbox should be queried for known root IDs after inventory growth"
|
||||
);
|
||||
let sync_log = std::fs::read_to_string(syncing.log_path()).expect("read syncing relay log");
|
||||
assert!(
|
||||
sync_log.lines().any(|line| {
|
||||
line.contains("Reconciled proactive participant mailbox inventory")
|
||||
&& line.contains("repository_mailbox_author_count=1")
|
||||
}),
|
||||
"repository mailbox ownership should remain observable without public-key labels"
|
||||
);
|
||||
assert!(
|
||||
!sync_log
|
||||
.lines()
|
||||
.any(|line| { line.contains("Starting fresh_start") && line.contains(mailbox.url()) }),
|
||||
"an owner mailbox must remain history-only rather than becoming a repository source"
|
||||
);
|
||||
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user