mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
feat(sync): follow addressable descendant references
Repository descendants can refer to a direct thread member by its NIP-01 address rather than its event ID. The existing non-recursive frontier retained only event IDs, so a/A tags and coordinate-valued q tags could fall outside both live and historic coverage. Represent the frontier as event IDs plus coordinates derived from valid replaceable and addressable direct members. Feed both sets through the existing live admission, five-second fallback rotation, and ordinary historic sync paths; malformed addressable events without d do not contribute a coordinate. This preserves the one-generation boundary and existing connection-budget policy. Recursive discovery and connection sharding remain deliberately excluded. Validated with 682 passing library tests, focused coordinate/filter unit tests, and all three descendant integration scenarios passing together. The constrained scenarios recover event-ID and a/A/q coordinate-only descendants through queued EOSE-closing history; their advertised capacity is fixed at four so minimum-churn packing cannot legitimately admit auxiliary live coverage.
This commit is contained in:
+4
-2
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
+144
-32
@@ -699,9 +699,15 @@ pub enum PendingBatchPurpose {
|
||||
Descendants,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||
struct DescendantFrontier {
|
||||
event_ids: HashSet<EventId>,
|
||||
coordinates: HashSet<String>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DescendantSyncRotation {
|
||||
frontier: HashSet<EventId>,
|
||||
frontier: DescendantFrontier,
|
||||
filters: Vec<DescendantFilterCursor>,
|
||||
next_filter: usize,
|
||||
in_flight: Option<DescendantFilterInFlight>,
|
||||
@@ -722,14 +728,14 @@ struct DescendantFilterInFlight {
|
||||
|
||||
#[derive(Debug)]
|
||||
struct DescendantLiveCoverage {
|
||||
frontier: HashSet<EventId>,
|
||||
frontier: DescendantFrontier,
|
||||
subscription_ids: Vec<SubscriptionId>,
|
||||
}
|
||||
|
||||
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<EventId>) {
|
||||
fn refresh(&mut self, frontier: DescendantFrontier) {
|
||||
let previous: HashMap<String, Option<Timestamp>> = 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<String> {
|
||||
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<Timestamp>,
|
||||
) -> Vec<Filter> {
|
||||
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<EventId>) -> HashSet<EventId> {
|
||||
let mut members = HashSet::new();
|
||||
async fn direct_thread_members(&self, root_events: &HashSet<EventId>) -> 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<EventId>,
|
||||
) {
|
||||
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<Timestamp>,
|
||||
) -> Option<u64> {
|
||||
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::<Vec<_>>();
|
||||
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");
|
||||
|
||||
@@ -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::<Vec<_>>();
|
||||
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::<Vec<_>>();
|
||||
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;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user