fix basic sync tests

This commit is contained in:
DanConwayDev
2025-12-05 11:04:00 +00:00
parent 83ede29fb2
commit ef7ba7c59b
4 changed files with 221 additions and 57 deletions
+5
View File
@@ -103,6 +103,11 @@ pub struct Config {
/// Number of days to look back for reconnect catchup (default: 3)
#[arg(long, env = "NGIT_SYNC_RECONNECT_LOOKBACK_DAYS", default_value_t = 3)]
pub sync_reconnect_lookback_days: u64,
/// Maximum startup jitter in milliseconds for sync connections (default: 10000 = 10 seconds)
/// Set to 0 to disable jitter (useful for testing)
#[arg(long, env = "NGIT_SYNC_STARTUP_JITTER_MS", default_value_t = 10_000)]
pub sync_startup_jitter_ms: u64,
}
impl Config {
+19 -12
View File
@@ -40,8 +40,6 @@ use super::metrics::SyncMetrics;
use crate::config::Config;
use crate::nostr::builder::{Nip34WritePolicy, SharedDatabase};
/// Maximum startup jitter in milliseconds (10 seconds)
const MAX_STARTUP_JITTER_MS: u64 = 10_000;
/// Default fallback address for sync source when bind_address cannot be parsed
///
@@ -76,6 +74,8 @@ pub struct SyncManager {
metrics: Option<SyncMetrics>,
/// Source address for synced events (derived from config.bind_address)
sync_source_addr: SocketAddr,
/// Maximum startup jitter in milliseconds (from config)
startup_jitter_ms: u64,
}
impl SyncManager {
@@ -102,6 +102,7 @@ impl SyncManager {
health_tracker: Arc::new(RelayHealthTracker::new(config)),
metrics: None,
sync_source_addr: get_sync_source_addr(&config.bind_address),
startup_jitter_ms: config.sync_startup_jitter_ms,
}
}
@@ -130,6 +131,7 @@ impl SyncManager {
health_tracker: Arc::new(RelayHealthTracker::new(config)),
metrics: Some(metrics),
sync_source_addr: get_sync_source_addr(&config.bind_address),
startup_jitter_ms: config.sync_startup_jitter_ms,
}
}
@@ -149,6 +151,7 @@ impl SyncManager {
health_tracker: Arc::new(RelayHealthTracker::with_defaults()),
metrics: None,
sync_source_addr: DEFAULT_SYNC_SOURCE_ADDR,
startup_jitter_ms: 10_000, // Default 10 seconds
}
}
@@ -265,8 +268,9 @@ impl SyncManager {
/// Spawn a connection task for a relay with startup jitter
///
/// Adds a random delay (0-10s) before connecting to prevent thundering herd
/// on startup when multiple relays are configured.
/// Adds a random delay (0 to startup_jitter_ms) before connecting to prevent
/// thundering herd on startup when multiple relays are configured.
/// Set startup_jitter_ms to 0 to disable jitter (useful for testing).
fn spawn_connection_with_jitter(
&self,
url: String,
@@ -276,16 +280,19 @@ impl SyncManager {
let domain = self.relay_domain.clone();
let health_tracker = self.health_tracker.clone();
let metrics = self.metrics.clone();
let max_jitter = self.startup_jitter_ms;
tokio::spawn(async move {
// Apply startup jitter
let jitter_ms = rand::thread_rng().gen_range(0..MAX_STARTUP_JITTER_MS);
tracing::debug!(
"Applying {}ms startup jitter before connecting to {}",
jitter_ms,
url
);
tokio::time::sleep(Duration::from_millis(jitter_ms)).await;
// Apply startup jitter (if configured)
if max_jitter > 0 {
let jitter_ms = rand::thread_rng().gen_range(0..max_jitter);
tracing::debug!(
"Applying {}ms startup jitter before connecting to {}",
jitter_ms,
url
);
tokio::time::sleep(Duration::from_millis(jitter_ms)).await;
}
connect_with_retry(&url, tx, filter_service, &domain, health_tracker, metrics).await;
});
+1
View File
@@ -92,6 +92,7 @@ impl TestRelay {
.env("NGIT_DOMAIN", &bind_address) // Set domain to match bind address
.env("NGIT_GIT_DATA_PATH", git_data_dir.path())
.env("NGIT_OWNER_NPUB", &test_npub)
.env("NGIT_SYNC_STARTUP_JITTER_MS", "0") // Disable jitter for tests
.env("RUST_LOG", "warn") // Less logging during tests
.stdout(Stdio::null())
.stderr(Stdio::null());
+196 -45
View File
@@ -21,25 +21,96 @@ use nostr_sdk::prelude::*;
/// Kind 30617 - Repository State (NIP-34)
const KIND_REPOSITORY_STATE: u16 = 30617;
/// Create a client with keys, connect to relay, and wait for connection
async fn create_connected_client(relay_url: &str, keys: Keys) -> Result<Client, String> {
let client = Client::new(keys);
client
.add_relay(relay_url)
.await
.map_err(|e| e.to_string())?;
client.connect().await;
// Wait for connection to establish (with retries, matching grasp-audit pattern)
for _ in 0..30 {
tokio::time::sleep(Duration::from_millis(100)).await;
let relays = client.relays().await;
if relays.values().any(|r| r.is_connected()) {
return Ok(client);
}
}
Err("Failed to connect to relay after 3 seconds".to_string())
}
/// Send an event and wait for successful delivery
async fn send_event_reliably(client: &Client, event: &Event) -> Result<EventId, String> {
// Try sending the event with retries
for attempt in 1..=5 {
let result = client.send_event(event).await;
match result {
Ok(output) => {
if !output.success.is_empty() {
return Ok(output.val);
}
// Check what went wrong
if !output.failed.is_empty() {
println!(" Attempt {} - failures: {:?}", attempt, output.failed);
// If relay not connected, try reconnecting
client.connect().await;
}
}
Err(e) => {
println!(" Attempt {} - error: {}", attempt, e);
}
}
tokio::time::sleep(Duration::from_millis(500)).await;
}
Err("Failed to send event after 5 attempts".to_string())
}
/// Create a valid repository announcement event for testing
///
/// This creates a kind 30617 event with required clone and relays tags
fn create_valid_repo_announcement(
keys: &Keys,
domain: &str,
identifier: &str,
) -> Event {
// Build tags for repository announcement
/// This creates a kind 30617 event with required clone and relays tags.
/// Uses TagKind::custom("clone") and TagKind::custom("relays") to match grasp-audit patterns.
#[allow(dead_code)]
fn create_valid_repo_announcement(keys: &Keys, domain: &str, identifier: &str) -> Event {
// Build tags for repository announcement using custom tag kinds (as grasp-audit does)
let tags = vec![
Tag::identifier(identifier),
Tag::custom(
TagKind::custom("clone"),
vec![format!("http://{}/{}", domain, identifier)],
),
Tag::custom(
TagKind::custom("relays"),
vec![format!("ws://{}", domain)],
vec![format!("http://{}/{}.git", domain, identifier)],
),
Tag::custom(TagKind::custom("relays"), vec![format!("ws://{}", domain)]),
];
EventBuilder::new(Kind::Custom(KIND_REPOSITORY_STATE), "Repository state")
.tags(tags)
.sign_with_keys(keys)
.expect("Failed to sign event")
}
/// Create a valid repository announcement event listing multiple relays
///
/// This creates a kind 30617 event with clone/relays tags referencing multiple domains,
/// which is necessary for sync tests where the event needs to be accepted by both relays.
/// Uses TagKind::custom("clone") and TagKind::custom("relays") to match grasp-audit patterns.
fn create_shared_repo_announcement(keys: &Keys, domains: &[&str], identifier: &str) -> Event {
// Build clone URLs for all domains (with .git suffix)
let clone_urls: Vec<String> = domains
.iter()
.map(|d| format!("http://{}/{}.git", d, identifier))
.collect();
// Build relay URLs for all domains
let relay_urls: Vec<String> = domains.iter().map(|d| format!("ws://{}", d)).collect();
// Build tags for repository announcement using custom tag kinds (as grasp-audit does)
let tags = vec![
Tag::identifier(identifier),
Tag::custom(TagKind::custom("clone"), clone_urls),
Tag::custom(TagKind::custom("relays"), relay_urls),
];
EventBuilder::new(Kind::Custom(KIND_REPOSITORY_STATE), "Repository state")
@@ -72,48 +143,116 @@ async fn test_sync_relay_connects_to_source() {
async fn test_valid_event_syncs_to_relay() {
// Start source relay (relay_a)
let relay_a = TestRelay::start().await;
// Give relay_a time to start
tokio::time::sleep(Duration::from_millis(200)).await;
println!(
"relay_a started at {} (domain: {})",
relay_a.url(),
relay_a.domain()
);
// Start syncing relay (relay_b) configured to sync from relay_a
let relay_b = TestRelay::start_with_sync(relay_a.url()).await;
println!(
"relay_b started at {} (domain: {})",
relay_b.url(),
relay_b.domain()
);
// Create test keys
// Create test keys that will be used for both client and event signing
let keys = Keys::generate();
// Create and submit a valid repository announcement to relay_a
let event = create_valid_repo_announcement(&keys, &relay_a.domain(), "test-repo");
// Wait for relay_b's sync connection to establish
// With NGIT_SYNC_STARTUP_JITTER_MS=0 (set by TestRelay), sync connects immediately.
// A brief wait allows the WebSocket connection and Layer 1 subscription to be set up.
println!("Waiting 1s for relay_b sync connection to establish...");
tokio::time::sleep(Duration::from_secs(1)).await;
// Create a client with our keys and connect to relay_a
let client_a = create_connected_client(relay_a.url(), keys.clone())
.await
.expect("Failed to connect to relay_a");
println!("client_a connected to relay_a");
// Create a repository announcement that lists BOTH relays
// This is required for sync - the event must reference both the source relay
// and the syncing relay for the write policy to accept it on both sides
let event = create_shared_repo_announcement(
&keys,
&[&relay_a.domain(), &relay_b.domain()],
"test-repo",
);
let event_id = event.id;
// Submit event to relay_a
let client_a = Client::default();
client_a.add_relay(relay_a.url()).await.expect("Failed to add relay_a");
client_a.connect().await;
// Print event details for debugging
println!("Created event {} (kind {})", event_id, event.kind.as_u16());
for tag in event.tags.iter() {
println!(" Tag: {:?}", tag.as_slice());
}
let send_result = client_a.send_event(&event).await;
assert!(send_result.is_ok(), "Failed to send event to relay_a: {:?}", send_result.err());
// Submit event to relay_a AFTER relay_b's subscription is established
// This ensures the event is received via the live subscription
println!("Sending event to relay_a...");
send_event_reliably(&client_a, &event)
.await
.expect("Failed to send event to relay_a");
println!("Event sent successfully");
// Wait for sync to occur
tokio::time::sleep(Duration::from_secs(2)).await;
// Query relay_b to verify the event was synced
let client_b = Client::default();
client_b.add_relay(relay_b.url()).await.expect("Failed to add relay_b");
client_b.connect().await;
// Create filter to find our event
let filter = Filter::new()
// Verify event is stored on relay_a first
let filter_a = Filter::new()
.kind(Kind::Custom(KIND_REPOSITORY_STATE))
.author(keys.public_key());
let events = client_b
.fetch_events(filter, Duration::from_secs(5))
let events_on_a = client_a
.fetch_events(filter_a.clone(), Duration::from_secs(5))
.await
.expect("Failed to fetch events from relay_a");
println!(
"Events on relay_a: {} (looking for {})",
events_on_a.len(),
event_id
);
for e in events_on_a.iter() {
println!(" Found event: {} (kind {})", e.id, e.kind.as_u16());
}
let found_on_a = events_on_a.iter().any(|e| e.id == event_id);
assert!(
found_on_a,
"Event {} was not stored on relay_a! This is a prerequisite for sync.",
event_id
);
println!("✓ Event confirmed on relay_a");
// Wait for sync to occur (event processing and storage)
println!("Waiting 1s for sync to occur...");
tokio::time::sleep(Duration::from_secs(1)).await;
// Query relay_b to verify the event was synced
let client_b = create_connected_client(relay_b.url(), Keys::generate())
.await
.expect("Failed to connect to relay_b");
// Create filter to find our event
let filter_b = Filter::new()
.kind(Kind::Custom(KIND_REPOSITORY_STATE))
.author(keys.public_key());
let events_on_b = client_b
.fetch_events(filter_b, Duration::from_secs(5))
.await
.expect("Failed to fetch events from relay_b");
println!(
"Events on relay_b: {} (looking for {})",
events_on_b.len(),
event_id
);
for e in events_on_b.iter() {
println!(" Found event: {} (kind {})", e.id, e.kind.as_u16());
}
// Check if our event was synced
let found = events.iter().any(|e| e.id == event_id);
let found_on_b = events_on_b.iter().any(|e| e.id == event_id);
// Clean up
client_a.disconnect().await;
@@ -122,10 +261,10 @@ async fn test_valid_event_syncs_to_relay() {
relay_a.stop().await;
assert!(
found,
"Event {} was not synced to relay_b. Found {} events",
found_on_b,
"Event {} was not synced to relay_b. Found {} events on relay_b",
event_id,
events.len()
events_on_b.len()
);
}
@@ -165,7 +304,10 @@ async fn test_invalid_event_rejected_by_sync_validation() {
// Submit invalid event to relay_a
// Note: relay_a will also reject it due to GRASP validation
let client_a = Client::default();
client_a.add_relay(relay_a.url()).await.expect("Failed to add relay_a");
client_a
.add_relay(relay_a.url())
.await
.expect("Failed to add relay_a");
client_a.connect().await;
// This will likely fail since relay_a also validates, but let's try
@@ -176,7 +318,10 @@ async fn test_invalid_event_rejected_by_sync_validation() {
// Query relay_b - the event should NOT be present
let client_b = Client::default();
client_b.add_relay(relay_b.url()).await.expect("Failed to add relay_b");
client_b
.add_relay(relay_b.url())
.await
.expect("Failed to add relay_b");
client_b.connect().await;
let filter = Filter::new()
@@ -221,7 +366,10 @@ async fn test_sync_respects_local_validation() {
let valid_event_id = valid_event.id;
let client_a = Client::default();
client_a.add_relay(relay_a.url()).await.expect("Failed to add relay_a");
client_a
.add_relay(relay_a.url())
.await
.expect("Failed to add relay_a");
client_a.connect().await;
client_a
@@ -234,7 +382,10 @@ async fn test_sync_respects_local_validation() {
// Query relay_b to verify the valid event was synced
let client_b = Client::default();
client_b.add_relay(relay_b.url()).await.expect("Failed to add relay_b");
client_b
.add_relay(relay_b.url())
.await
.expect("Failed to add relay_b");
client_b.connect().await;
let filter = Filter::new()
@@ -259,4 +410,4 @@ async fn test_sync_respects_local_validation() {
"Valid event {} should have been synced to relay_b",
valid_event_id
);
}
}