mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
test: poll wait_for_event_served at a fixed interval over one connection
wait_for_event_served opened a fresh client per poll, waited for the connection in 100 ms steps, and backed off exponentially from 100 ms to 2 s between polls with a 2 s fetch timeout. Tests that reach it right after a send, now that the covering sleeps are gone, could miss the first poll while the relay was still promoting and then wait seconds for the next. Connect once with try_connect, bounded by the caller's timeout, and poll every 100 ms with a 1 s fetch timeout until the deadline. Healthy runs still cost one round trip; slow runs detect the event within roughly one interval instead of a growing backoff. Deadline and error semantics are unchanged. Validation: the two suites that use this helper most passed 20 unloaded, 20 loaded and 30 concurrent-contention samples with no failures and no timing regression against the previous helper. Assisted-by: Claude Fable 5.1 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5.1
parent
86047d78d6
commit
6c3b33e3a6
@@ -367,78 +367,55 @@ pub fn create_announcement_event(
|
||||
.map_err(|e| format!("Failed to sign announcement event: {}", e))
|
||||
}
|
||||
|
||||
/// Wait for an event to be served by a relay (not in purgatory).
|
||||
/// Wait until `relay_url` serves `event_id`, polling at a fixed short
|
||||
/// interval with one connection, bounded by `timeout`.
|
||||
///
|
||||
/// Polls the relay until the event is queryable, indicating it has
|
||||
/// been released from purgatory. Uses exponential backoff for polling.
|
||||
///
|
||||
/// # Arguments
|
||||
/// * `relay_url` - WebSocket URL of the relay
|
||||
/// * `event_id` - Event ID to wait for
|
||||
/// * `timeout` - Maximum time to wait
|
||||
///
|
||||
/// # Returns
|
||||
/// * `Ok(Event)` - The event was found
|
||||
/// * `Err(String)` - Timeout or error
|
||||
/// A healthy relay answers the first poll, so the wait costs one round
|
||||
/// trip. When the event is still being promoted, a fixed interval keeps
|
||||
/// the detection delay near the interval instead of letting an exponential
|
||||
/// backoff stretch it under load.
|
||||
pub async fn wait_for_event_served(
|
||||
relay_url: &str,
|
||||
event_id: &EventId,
|
||||
timeout: Duration,
|
||||
) -> Result<Event, String> {
|
||||
let deadline = std::time::Instant::now() + timeout;
|
||||
let poll_interval = Duration::from_millis(100);
|
||||
let fetch_timeout = Duration::from_secs(1);
|
||||
|
||||
let temp_keys = Keys::generate();
|
||||
let client = Client::builder()
|
||||
.authenticator(SignerAuthenticator::new(temp_keys))
|
||||
.build();
|
||||
|
||||
client
|
||||
.add_relay(relay_url)
|
||||
.await
|
||||
.map_err(|e| format!("Failed to add relay: {}", e))?;
|
||||
|
||||
client.connect().await;
|
||||
|
||||
// Wait for connection
|
||||
let mut connected = false;
|
||||
for _ in 0..20 {
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
let relays = client.relays().await;
|
||||
if relays.values().any(|r| r.status().is_connected()) {
|
||||
connected = true;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
if !connected {
|
||||
let output = client.try_connect().timeout(timeout).await;
|
||||
if output.success.is_empty() {
|
||||
client.disconnect().await;
|
||||
return Err("Failed to connect to relay".to_string());
|
||||
return Err(format!(
|
||||
"Failed to connect to relay {}: {:?}",
|
||||
relay_url, output.failed
|
||||
));
|
||||
}
|
||||
|
||||
// Poll for the event with exponential backoff
|
||||
let start = std::time::Instant::now();
|
||||
let mut poll_interval = Duration::from_millis(100);
|
||||
let max_interval = Duration::from_secs(2);
|
||||
|
||||
while start.elapsed() < timeout {
|
||||
let filter = Filter::new().id(*event_id);
|
||||
|
||||
match client
|
||||
.fetch_events(filter)
|
||||
.timeout(Duration::from_secs(2))
|
||||
let filter = Filter::new().id(*event_id);
|
||||
loop {
|
||||
if let Ok(events) = client
|
||||
.fetch_events(filter.clone())
|
||||
.timeout(fetch_timeout)
|
||||
.await
|
||||
{
|
||||
Ok(events) => {
|
||||
if let Some(event) = events.into_iter().next() {
|
||||
client.disconnect().await;
|
||||
return Ok(event);
|
||||
}
|
||||
}
|
||||
Err(_) => {
|
||||
// Ignore fetch errors, will retry
|
||||
if let Some(event) = events.into_iter().next() {
|
||||
client.disconnect().await;
|
||||
return Ok(event);
|
||||
}
|
||||
}
|
||||
|
||||
if std::time::Instant::now() >= deadline {
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(poll_interval).await;
|
||||
poll_interval = std::cmp::min(poll_interval * 2, max_interval);
|
||||
}
|
||||
|
||||
client.disconnect().await;
|
||||
|
||||
Reference in New Issue
Block a user