fix: enable Layer 3 sync by adding root events to pending queue

When root events (issues/patches) are received via self-subscription,
handle_root_event() was only updating the repo_sync_index directly.
This caused process_batch() to early-return when pending.is_empty(),
so Layer 3 filters for comments/replies were never created.

The fix adds root events to both:
1. repo_sync_index (for immediate availability)
2. pending queue (to trigger Layer 3 filter creation in next batch)

Critical: The pending entry must include relays from repo_sync_index
so derive_relay_targets() knows where to send Layer 3 subscriptions.

The Layer 3 test now verifies that events sent BEFORE the subscription
is established are still synced - proving subscriptions without 'since'
correctly fetch historical events.

Enabled 4 previously ignored Layer 3 tests:
- test_live_sync_layer3_events
- test_layer3_sync_with_lowercase_e_tag
- test_layer3_sync_with_uppercase_e_tag
- test_layer3_sync_with_q_tag
This commit is contained in:
DanConwayDev
2025-12-10 21:57:02 +00:00
parent 466a009d82
commit 46b306dcfa
3 changed files with 98 additions and 77 deletions
+25 -7
View File
@@ -167,13 +167,14 @@ impl SelfSubscriber {
}
} else {
// For root event kinds (1617, 1618, 1619, 1621),
// process them to update the RepoSyncIndex
// process them to update the RepoSyncIndex AND add to pending
// for Layer 3 filter creation
tracing::trace!(
kind = %event.kind,
event_id = %event.id,
"Received root event"
);
self.handle_root_event(&event).await;
self.handle_root_event(&event, pending).await;
}
LoopControl::Continue
}
@@ -375,8 +376,9 @@ impl SelfSubscriber {
/// Handle a root event (1617/1618/1619/1621)
///
/// Extracts the 'a' tag to find the repo addressable reference,
/// then updates the RepoSyncIndex with the event ID.
async fn handle_root_event(&self, event: &Event) {
/// then updates the RepoSyncIndex with the event ID AND adds to pending
/// so that Layer 3 filters will be created in the next batch.
async fn handle_root_event(&self, event: &Event, pending: &mut PendingUpdates) {
// Extract 'a' tag to find the repo addressable reference
let repo_a_tag = event.tags.iter().find(|tag| {
let tag_vec = tag.as_slice();
@@ -406,15 +408,31 @@ impl SelfSubscriber {
}
};
// Look up repo in repo_sync_index
// Look up repo in repo_sync_index - add root event directly and also to pending
let mut index = self.repo_sync_index.write().await;
if let Some(repo_sync) = index.get_mut(&repo_ref) {
// Add event.id to root_events set
// Add event.id to root_events set in the index (immediate availability)
repo_sync.root_events.insert(event.id);
// Clone the relays before releasing the lock - Layer 3 filters need to be
// sent to the same relays as Layer 2 filters for this repo
let relays = repo_sync.relays.clone();
// Release lock before modifying pending
drop(index);
// Also add root event to pending - this ensures batch processing runs
// and creates Layer 3 filters for events referencing this root event.
// CRITICAL: Include relays so derive_relay_targets knows where to send filters!
let mut root_events = HashSet::new();
root_events.insert(event.id);
pending.add_repo(repo_ref.clone(), relays.clone(), root_events);
tracing::debug!(
event_id = %event.id,
repo_ref = %repo_ref,
"Added root event to repo"
relay_count = relays.len(),
"Added root event to index and pending for Layer 3 filter creation"
);
} else {
tracing::debug!(
+49 -49
View File
@@ -58,11 +58,8 @@ async fn test_live_sync_layer2_events() {
// 4. Create a repository announcement that lists BOTH relays
let repo_id = "test-repo-live-l2";
let announcement = create_repo_announcement(
&keys,
&[&relay_a.domain(), &relay_b.domain()],
repo_id,
);
let announcement =
create_repo_announcement(&keys, &[&relay_a.domain(), &relay_b.domain()], repo_id);
println!(
"Created announcement {} (kind {})",
@@ -123,7 +120,7 @@ async fn test_live_sync_layer2_events() {
.id(issue_id);
let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(5)).await;
println!("Issue {} synced to relay_b: {}", issue_id, synced);
// 10. Cleanup
@@ -150,15 +147,7 @@ async fn test_live_sync_layer2_events() {
/// 6. Verify comment syncs to relay_b within 5 seconds
/// 7. Verify comment has correct 'E' tag reference
///
/// # Note
/// This test is currently ignored because Layer 3 (comment) sync is tracked separately
/// and may not be fully implemented yet. See discovery.rs for context:
/// > "Note: Layer 3 (comments on issues) sync is tracked separately and may
/// > be implemented in future phases."
///
/// TODO: Enable this test when Layer 3 sync is implemented.
#[tokio::test]
#[ignore = "Layer 3 sync not yet implemented - comments don't sync via discovery"]
async fn test_live_sync_layer3_events() {
// 1. Start relays
let relay_a = TestRelay::start().await;
@@ -179,11 +168,8 @@ async fn test_live_sync_layer3_events() {
// 2. Create and send repository announcement to both relays
let repo_id = "test-repo-live-l3";
let announcement = create_repo_announcement(
&keys,
&[&relay_a.domain(), &relay_b.domain()],
repo_id,
);
let announcement =
create_repo_announcement(&keys, &[&relay_a.domain(), &relay_b.domain()], repo_id);
let client_a = TestClient::new(relay_a.url(), keys.clone())
.await
@@ -220,16 +206,8 @@ async fn test_live_sync_layer3_events() {
.expect("Failed to send issue");
println!("Layer 2 issue {} sent to relay_a", issue_id);
// 5. Wait for issue to sync to relay_b
tokio::time::sleep(Duration::from_secs(2)).await;
let issue_filter = Filter::new()
.kind(Kind::Custom(KIND_ISSUE))
.id(issue_id);
let issue_synced = wait_for_event_on_relay(relay_b.url(), issue_filter, Duration::from_secs(3)).await;
println!("Issue synced to relay_b: {}", issue_synced);
// 6. Create and send Layer 3 comment (kind 1111, uppercase E tag)
// 5. Create and send Layer 3 comment IMMEDIATELY (before waiting for sync)
// This tests that subscriptions without 'since' will catch pre-existing events
let comment = build_layer3_comment_with_uppercase_e_tag(
&keys,
&issue_id,
@@ -238,7 +216,11 @@ async fn test_live_sync_layer3_events() {
.expect("Failed to create comment");
let comment_id = comment.id;
println!("Created comment {} (kind {})", comment_id, comment.kind.as_u16());
println!(
"Created comment {} (kind {})",
comment_id,
comment.kind.as_u16()
);
for tag in comment.tags.iter() {
println!(" Tag: {:?}", tag.as_slice());
}
@@ -247,7 +229,15 @@ async fn test_live_sync_layer3_events() {
.send_event(&comment)
.await
.expect("Failed to send comment");
println!("Layer 3 comment {} sent to relay_a", comment_id);
println!("Layer 3 comment {} sent to relay_a BEFORE Layer 3 subscription established", comment_id);
// 6. Now wait for issue to sync to relay_b (this triggers Layer 3 filter creation)
tokio::time::sleep(Duration::from_secs(2)).await;
let issue_filter = Filter::new().kind(Kind::Custom(KIND_ISSUE)).id(issue_id);
let issue_synced =
wait_for_event_on_relay(relay_b.url(), issue_filter, Duration::from_secs(3)).await;
println!("Issue synced to relay_b: {}", issue_synced);
client_a.disconnect().await;
client_b.disconnect().await;
@@ -258,8 +248,12 @@ async fn test_live_sync_layer3_events() {
.author(keys.public_key())
.id(comment_id);
let comment_synced = wait_for_event_on_relay(relay_b.url(), comment_filter, Duration::from_secs(5)).await;
println!("Comment {} synced to relay_b: {}", comment_id, comment_synced);
let comment_synced =
wait_for_event_on_relay(relay_b.url(), comment_filter, Duration::from_secs(5)).await;
println!(
"Comment {} synced to relay_b: {}",
comment_id, comment_synced
);
// 8. Verify the comment has correct 'E' tag reference
let mut has_correct_ref = false;
@@ -269,12 +263,15 @@ async fn test_live_sync_layer3_events() {
if client.add_relay(relay_b.url()).await.is_ok() {
client.connect().await;
tokio::time::sleep(Duration::from_millis(500)).await;
let fetch_filter = Filter::new()
.kind(Kind::Custom(KIND_COMMENT))
.id(comment_id);
if let Ok(events) = client.fetch_events(fetch_filter, Duration::from_secs(2)).await {
if let Ok(events) = client
.fetch_events(fetch_filter, Duration::from_secs(2))
.await
{
if let Some(event) = events.first() {
// Check for 'E' tag with parent event ID
for tag in event.tags.iter() {
@@ -347,11 +344,8 @@ async fn test_live_sync_event_ordering() {
// 2. Create and send repository announcement to both relays
let repo_id = "test-repo-ordering";
let announcement = create_repo_announcement(
&keys,
&[&relay_a.domain(), &relay_b.domain()],
repo_id,
);
let announcement =
create_repo_announcement(&keys, &[&relay_a.domain(), &relay_b.domain()], repo_id);
let client_a = TestClient::new(relay_a.url(), keys.clone())
.await
@@ -387,12 +381,15 @@ async fn test_live_sync_event_ordering() {
&format!("Ordering Test Issue {}", i),
)
.expect("Failed to create issue");
// Store the created_at timestamp for ordering verification
expected_order_timestamps.push(issue.created_at);
issue_ids.push(issue.id);
println!("Created issue {} at timestamp {}", issue.id, issue.created_at);
println!(
"Created issue {} at timestamp {}",
issue.id, issue.created_at
);
client_a
.send_event(&issue)
@@ -412,7 +409,7 @@ async fn test_live_sync_event_ordering() {
// 6. Fetch all events from relay_b
let temp_keys = Keys::generate();
let client = Client::new(temp_keys);
let events_found: Vec<Event>;
if client.add_relay(relay_b.url()).await.is_ok() {
client.connect().await;
@@ -446,7 +443,11 @@ async fn test_live_sync_event_ordering() {
.filter(|e| issue_ids.contains(&e.id))
.collect();
println!("Found {} test events (out of {} total)", test_events.len(), events_found.len());
println!(
"Found {} test events (out of {} total)",
test_events.len(),
events_found.len()
);
// 8. Check ordering by created_at timestamp
let mut ordered_correctly = true;
@@ -470,8 +471,7 @@ async fn test_live_sync_event_ordering() {
ordered_correctly = false;
println!(
"Order violation: {} ({}) > {} ({})",
window[0].id, window[0].created_at,
window[1].id, window[1].created_at
window[0].id, window[0].created_at, window[1].id, window[1].created_at
);
}
}
@@ -492,4 +492,4 @@ async fn test_live_sync_event_ordering() {
ordered_correctly,
"Events should be ordered by created_at timestamp"
);
}
}
+24 -21
View File
@@ -329,14 +329,7 @@ async fn test_layer2_sync_with_q_tag() {
/// tags sync correctly between relays when referencing a Layer 2 event.
///
/// The lowercase 'e' tag is the standard NIP-01 way to reference events by ID.
///
/// # Note
/// This test is currently ignored because Layer 3 sync is not yet implemented.
/// The test logic is complete and will work when Layer 3 sync is enabled.
///
/// TODO: Enable this test when Layer 3 sync is implemented.
#[tokio::test]
#[ignore = "Layer 3 sync not yet implemented - comments don't sync via discovery"]
async fn test_layer3_sync_with_lowercase_e_tag() {
// 1. Start relays
let relay_a = TestRelay::start().await;
@@ -406,6 +399,14 @@ async fn test_layer3_sync_with_lowercase_e_tag() {
println!("Issue synced to relay_b: {}", issue_synced);
assert!(issue_synced, "Layer 2 issue should sync first");
// Wait for Layer 3 subscriptions to be established
// After issue syncs, relay_b's SelfSubscriber needs time to:
// 1. Receive the synced issue via notify_event broadcast
// 2. Batch timer to tick (up to 200ms in tests)
// 3. Process batch and create Layer 3 filters
// 4. Subscribe to relay_a with Layer 3 filters
tokio::time::sleep(Duration::from_millis(500)).await;
// 6. Create and send Layer 3 reply with lowercase 'e' tag (kind 1)
let reply = build_layer3_reply_with_e_tag(&keys, &issue_id, "Reply with lowercase e tag")
.expect("Failed to create reply");
@@ -451,14 +452,7 @@ async fn test_layer3_sync_with_lowercase_e_tag() {
/// tags sync correctly between relays when referencing a Layer 2 event.
///
/// The uppercase 'E' tag is used in NIP-22 for comment events.
///
/// # Note
/// This test is currently ignored because Layer 3 sync is not yet implemented.
/// The test logic is complete and will work when Layer 3 sync is enabled.
///
/// TODO: Enable this test when Layer 3 sync is implemented.
#[tokio::test]
#[ignore = "Layer 3 sync not yet implemented - comments don't sync via discovery"]
async fn test_layer3_sync_with_uppercase_e_tag() {
// 1. Start relays
let relay_a = TestRelay::start().await;
@@ -528,6 +522,14 @@ async fn test_layer3_sync_with_uppercase_e_tag() {
println!("Issue synced to relay_b: {}", issue_synced);
assert!(issue_synced, "Layer 2 issue should sync first");
// Wait for Layer 3 subscriptions to be established
// After issue syncs, relay_b's SelfSubscriber needs time to:
// 1. Receive the synced issue via notify_event broadcast
// 2. Batch timer to tick (up to 200ms in tests)
// 3. Process batch and create Layer 3 filters
// 4. Subscribe to relay_a with Layer 3 filters
tokio::time::sleep(Duration::from_millis(500)).await;
// 6. Create and send Layer 3 comment with uppercase 'E' tag (kind 1111)
let comment = build_layer3_comment_with_uppercase_e_tag(&keys, &issue_id, "Comment with uppercase E tag")
.expect("Failed to create comment");
@@ -573,14 +575,7 @@ async fn test_layer3_sync_with_uppercase_e_tag() {
/// tags sync correctly between relays when referencing a Layer 2 event.
///
/// The 'q' tag is used in NIP-18 for quotes/reposts.
///
/// # Note
/// This test is currently ignored because Layer 3 sync is not yet implemented.
/// The test logic is complete and will work when Layer 3 sync is enabled.
///
/// TODO: Enable this test when Layer 3 sync is implemented.
#[tokio::test]
#[ignore = "Layer 3 sync not yet implemented - comments don't sync via discovery"]
async fn test_layer3_sync_with_q_tag() {
// 1. Start relays
let relay_a = TestRelay::start().await;
@@ -650,6 +645,14 @@ async fn test_layer3_sync_with_q_tag() {
println!("Issue synced to relay_b: {}", issue_synced);
assert!(issue_synced, "Layer 2 issue should sync first");
// Wait for Layer 3 subscriptions to be established
// After issue syncs, relay_b's SelfSubscriber needs time to:
// 1. Receive the synced issue via notify_event broadcast
// 2. Batch timer to tick (up to 200ms in tests)
// 3. Process batch and create Layer 3 filters
// 4. Subscribe to relay_a with Layer 3 filters
tokio::time::sleep(Duration::from_millis(500)).await;
// 6. Create and send Layer 3 quote with 'q' tag (kind 1)
let quote = build_layer3_quote_with_q_tag(&keys, &issue_id, "Quote with q tag")
.expect("Failed to create quote");