diff --git a/CHANGELOG.md b/CHANGELOG.md index 203d080..4f5c655 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,8 +13,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 but omit the repository and root-event tags. When the per-connection ledger can retain complete descendant coverage while preserving control and historic capacity, those filters stay live; otherwise bounded EOSE-closing - history provides eventual coverage. The frontier is deliberately - non-recursive and connection sharding remains unnecessary. + history provides eventual coverage. Direct replaceable and addressable + members contribute both their event ID and NIP-01 coordinate, covering + descendants which use `a`, `A`, or coordinate-valued `q` tags. The frontier + is deliberately non-recursive and connection sharding remains unnecessary. ### Changed diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index bd8f973..7d8b66f 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -317,7 +317,10 @@ slots. Four consumers share it, in priority order: its complete separately grouped set fits after core coverage while still preserving the margin and a transient slot. Otherwise each constrained relay advances one cursor-overlapped REQ+EOSE filter per five-second tick - through the same transient queue. + through the same transient queue. Direct thread members contribute their + event IDs; replaceable and addressable members also contribute their NIP-01 + coordinates. The resulting `e`/`E`/`q` and `a`/`A`/coordinate-`q` filters + cover one descendant generation without recursively expanding the frontier. NIP-11 `max_subscriptions` sets B for each new connection session; when it is absent B falls back to 20. Advertised values below that floor are honoured diff --git a/src/sync/mod.rs b/src/sync/mod.rs index d3cc367..c449236 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -699,9 +699,15 @@ pub enum PendingBatchPurpose { Descendants, } +#[derive(Debug, Clone, Default, PartialEq, Eq)] +struct DescendantFrontier { + event_ids: HashSet, + coordinates: HashSet, +} + #[derive(Debug)] struct DescendantSyncRotation { - frontier: HashSet, + frontier: DescendantFrontier, filters: Vec, next_filter: usize, in_flight: Option, @@ -722,14 +728,14 @@ struct DescendantFilterInFlight { #[derive(Debug)] struct DescendantLiveCoverage { - frontier: HashSet, + frontier: DescendantFrontier, subscription_ids: Vec, } impl Default for DescendantSyncRotation { fn default() -> Self { Self { - frontier: HashSet::new(), + frontier: DescendantFrontier::default(), filters: Vec::new(), next_filter: 0, in_flight: None, @@ -738,20 +744,20 @@ impl Default for DescendantSyncRotation { } impl DescendantSyncRotation { - fn refresh(&mut self, members: HashSet) { + fn refresh(&mut self, frontier: DescendantFrontier) { let previous: HashMap> = self .filters .drain(..) .map(|cursor| (cursor.filter.as_json(), cursor.last_successful_until)) .collect(); - self.frontier = members.clone(); - self.filters = filters::tagged_one_of_our_root_event_filters(&members, None) + self.filters = descendant_frontier_filters(&frontier, None) .into_iter() .map(|filter| DescendantFilterCursor { last_successful_until: previous.get(&filter.as_json()).copied().flatten(), filter, }) .collect(); + self.frontier = frontier; self.next_filter = 0; self.in_flight = None; } @@ -794,6 +800,40 @@ impl DescendantSyncRotation { } } +fn descendant_event_coordinate(event: &Event) -> Option { + if !(event.kind.is_replaceable() || event.kind.is_addressable()) { + return None; + } + let identifier = if event.kind.is_addressable() { + event + .tags + .iter() + .find(|tag| tag.kind() == "d") + .and_then(|tag| tag.content())? + } else { + "" + }; + Some(format!( + "{}:{}:{}", + event.kind.as_u16(), + event.pubkey.to_hex(), + identifier + )) +} + +fn descendant_frontier_filters( + frontier: &DescendantFrontier, + since: Option, +) -> Vec { + let mut filters = + filters::tagged_one_of_our_root_event_filters(&frontier.event_ids, since); + filters.extend(filters::tagged_one_of_our_repo_event_filters( + &frontier.coordinates, + since, + )); + filters +} + /// Items included in a pending batch #[derive(Debug, Clone, Default)] pub struct PendingItems { @@ -3681,17 +3721,20 @@ impl SyncManager { /// Only events that directly reference a root are admitted to this /// frontier. Events fetched by the descendant rotation are deliberately /// not fed back into it, keeping this feature non-recursive. - async fn direct_thread_members(&self, root_events: &HashSet) -> HashSet { - let mut members = HashSet::new(); + async fn direct_thread_members(&self, root_events: &HashSet) -> DescendantFrontier { + let mut frontier = DescendantFrontier::default(); for filter in filters::tagged_one_of_our_root_event_filters(root_events, None) { match self.database.query(filter).await { Ok(events) => { - members.extend( - events - .iter() - .map(|event| event.id) - .filter(|event_id| !root_events.contains(event_id)), - ); + for event in events + .iter() + .filter(|event| !root_events.contains(&event.id)) + { + frontier.event_ids.insert(event.id); + if let Some(coordinate) = descendant_event_coordinate(event) { + frontier.coordinates.insert(coordinate); + } + } } Err(error) => { tracing::warn!( @@ -3699,11 +3742,11 @@ impl SyncManager { root_event_count = root_events.len(), "Failed to derive direct repository thread members" ); - return HashSet::new(); + return DescendantFrontier::default(); } } } - members + frontier } async fn close_descendant_live_coverage( @@ -3744,8 +3787,8 @@ impl SyncManager { relay_url: &str, root_events: &HashSet, ) { - let members = self.direct_thread_members(root_events).await; - if members.is_empty() { + let frontier = self.direct_thread_members(root_events).await; + if frontier.event_ids.is_empty() { let _ = self .close_descendant_live_coverage(relay_url, "frontier became empty") .await; @@ -3755,7 +3798,7 @@ impl SyncManager { if self .descendant_live_coverage .get(relay_url) - .is_some_and(|coverage| coverage.frontier == members) + .is_some_and(|coverage| coverage.frontier == frontier) { return; } @@ -3772,9 +3815,8 @@ impl SyncManager { .as_secs() .saturating_sub(DESCENDANT_FALLBACK_OVERLAP_SECS), ); - let historic_filters = filters::tagged_one_of_our_root_event_filters(&members, None); - let live_filters = - filters::tagged_one_of_our_root_event_filters(&members, Some(live_since)); + let historic_filters = descendant_frontier_filters(&frontier, None); + let live_filters = descendant_frontier_filters(&frontier, Some(live_since)); if let Some(connection) = self.connections.get(relay_url).cloned() { let groups = live_filter_groups(&live_filters, connection.max_filters_per_req()); if connection.can_admit_auxiliary_live_groups(&groups) { @@ -3784,16 +3826,20 @@ impl SyncManager { { Ok(subscription_ids) => { let filter_count = live_filters.len(); + let event_id_count = frontier.event_ids.len(); + let coordinate_count = frontier.coordinates.len(); self.descendant_live_coverage.insert( relay_url.to_string(), DescendantLiveCoverage { - frontier: members, + frontier, subscription_ids: subscription_ids.clone(), }, ); self.descendant_sync_rotations.remove(relay_url); tracing::info!( relay = %relay_url, + event_id_count, + coordinate_count, filter_count, subscription_count = subscription_ids.len(), "Installed auxiliary descendant live coverage" @@ -3834,8 +3880,8 @@ impl SyncManager { .descendant_sync_rotations .entry(relay_url.to_string()) .or_default(); - if rotation.frontier != members { - rotation.refresh(members); + if rotation.frontier != frontier { + rotation.refresh(frontier); } } @@ -6336,10 +6382,8 @@ impl SyncManager { since: Option, ) -> Option { if !items.root_events.is_empty() { - let members = self.direct_thread_members(&items.root_events).await; - filters.extend(filters::tagged_one_of_our_root_event_filters( - &members, None, - )); + let frontier = self.direct_thread_members(&items.root_events).await; + filters.extend(descendant_frontier_filters(&frontier, None)); } self.historic_sync_with_options( relay_url, @@ -6739,6 +6783,65 @@ impl SyncManager { mod tests { use super::*; + #[test] + fn descendant_frontier_derives_replaceable_and_addressable_coordinates() { + let keys = Keys::generate(); + let addressable = EventBuilder::new(Kind::Custom(30_023), "addressable") + .tag(Tag::custom("d", ["article-name"])) + .finalize(&keys) + .unwrap(); + let replaceable = EventBuilder::new(Kind::Custom(10_000), "replaceable") + .finalize(&keys) + .unwrap(); + let regular = EventBuilder::new(Kind::TextNote, "regular") + .finalize(&keys) + .unwrap(); + let malformed_addressable = EventBuilder::new(Kind::Custom(30_023), "missing d") + .finalize(&keys) + .unwrap(); + + assert_eq!( + descendant_event_coordinate(&addressable), + Some(format!("30023:{}:article-name", keys.public_key().to_hex())) + ); + assert_eq!( + descendant_event_coordinate(&replaceable), + Some(format!("10000:{}:", keys.public_key().to_hex())) + ); + assert_eq!(descendant_event_coordinate(®ular), None); + assert_eq!(descendant_event_coordinate(&malformed_addressable), None); + } + + #[test] + fn descendant_frontier_queries_event_and_coordinate_tag_variants() { + let since = Timestamp::from_secs(1234); + let frontier = DescendantFrontier { + event_ids: HashSet::from([EventId::from_byte_array([7; 32])]), + coordinates: HashSet::from([format!("30023:{}:article", "a".repeat(64))]), + }; + let filters = descendant_frontier_filters(&frontier, Some(since)); + + assert_eq!(filters.len(), 6, "e/E/q and a/A/q must all be covered"); + assert!(filters.iter().all(|filter| { + serde_json::to_value(filter).unwrap()["since"] == serde_json::json!(1234) + })); + let serialized = filters.iter().map(Filter::as_json).collect::>(); + for tag in ["#e", "#E", "#a", "#A"] { + assert!( + serialized.iter().any(|filter| filter.contains(tag)), + "missing {tag} descendant filter" + ); + } + assert_eq!( + serialized + .iter() + .filter(|filter| filter.contains("#q")) + .count(), + 2, + "event IDs and coordinates each require quote coverage" + ); + } + #[test] fn own_relay_targets_are_excluded_from_sync_actions() { assert!(is_own_sync_target("wss://gitnostr.com", "gitnostr.com")); @@ -6770,11 +6873,14 @@ mod tests { #[test] fn descendant_rotation_advances_cursor_only_after_successful_eose() { let member = EventId::from_byte_array([7; 32]); - let members = HashSet::from([member]); + let frontier = DescendantFrontier { + event_ids: HashSet::from([member]), + coordinates: HashSet::new(), + }; let now = Timestamp::from_secs(200_000); let mut rotation = DescendantSyncRotation::default(); - rotation.refresh(members); + rotation.refresh(frontier); let (filter_index, first, until) = rotation.next_request(now).unwrap(); assert!(serde_json::to_value(first).unwrap().get("since").is_none()); rotation.mark_started(41, filter_index, until); @@ -6810,7 +6916,13 @@ mod tests { EventId::from_byte_array(bytes) }) .collect(); - let filters = filters::tagged_one_of_our_root_event_filters(&members, None); + let filters = descendant_frontier_filters( + &DescendantFrontier { + event_ids: members, + coordinates: HashSet::new(), + }, + None, + ); let groups = live_filter_groups(&filters, MAX_FILTERS_PER_REQ); assert!(filters.len() > 3, "the frontier must span byte chunks"); diff --git a/tests/sync/descendant_sync.rs b/tests/sync/descendant_sync.rs index f5c21fe..0fc8d79 100644 --- a/tests/sync/descendant_sync.rs +++ b/tests/sync/descendant_sync.rs @@ -27,7 +27,7 @@ async fn wait_for_log(path: &std::path::Path, needle: &str, timeout: Duration) - /// pass. The recovered events do not recursively extend the frontier. #[tokio::test] async fn historic_sync_recovers_one_generation_of_parent_only_descendants() { - let source = TestRelay::start_with_relay_max_subscriptions(5).await; + let source = TestRelay::start_with_relay_max_subscriptions(4).await; let keys = Keys::generate(); let repo_id = "scheduled-descendant-history"; @@ -164,3 +164,84 @@ async fn descendant_live_coverage_is_preferred_when_capacity_remains() { syncing.stop().await; source.stop().await; } + +/// A direct addressable thread member extends the non-recursive frontier by +/// its coordinate as well as its event ID. Descendants which carry only an +/// address tag must therefore be recovered by a constrained relay's historic +/// fallback rotation. +#[tokio::test] +async fn historic_sync_recovers_address_tag_descendants() { + let source = TestRelay::start_with_relay_max_subscriptions(4).await; + let keys = Keys::generate(); + let repo_id = "addressable-descendant-history"; + let source_domains = [source.domain()]; + let source_refs = source_domains + .iter() + .map(String::as_str) + .collect::>(); + let (_announcement, _source_git) = + setup_announcement_on_relay(&source, &keys, &source_refs, repo_id).await; + let source_client = TestClient::new(source.url(), keys.clone()) + .await + .expect("connect to source relay"); + let issue = build_layer2_issue_event( + &keys, + &repo_coord(&keys, repo_id), + "Addressable descendant root", + ) + .expect("build issue"); + let direct_addressable = EventBuilder::new(Kind::Custom(30_023), "Direct article") + .tags([ + Tag::custom("d", ["direct-article"]), + Tag::custom("e", [issue.id.to_hex()]), + ]) + .finalize(&keys) + .expect("build direct addressable member"); + let coordinate = format!( + "30023:{}:direct-article", + direct_addressable.pubkey.to_hex() + ); + let address_only_children = ["a", "A", "q"].map(|tag| { + EventBuilder::new(Kind::Custom(1_111), format!("{tag}-tag child")) + .tag(Tag::custom(tag, [coordinate.clone()])) + .finalize(&keys) + .expect("build address-only child") + }); + source_client.send_event(&issue).await.unwrap(); + source_client.send_event(&direct_addressable).await.unwrap(); + for event in &address_only_children { + source_client.send_event(event).await.unwrap(); + } + + let syncing = TestRelay::start_with_sync(None).await; + let domains = [source.domain(), syncing.domain()]; + let domain_refs = domains.iter().map(String::as_str).collect::>(); + let (_target_announcement, _target_git) = + setup_announcement_on_relay(&syncing, &keys, &domain_refs, repo_id).await; + + for event in &address_only_children { + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(event.id), + Duration::from_secs(25), + ) + .await, + "historic fallback should recover the {} child", + event.content + ); + } + assert!( + wait_for_log( + &syncing.log_path(), + "Started queued descendant fallback query", + Duration::from_secs(5), + ) + .await, + "the constrained source should exercise coordinate fallback filters" + ); + + source_client.disconnect().await; + syncing.stop().await; + source.stop().await; +}