diff --git a/src/purgatory/persistence.rs b/src/purgatory/persistence.rs index 7fca2cf..67cd129 100644 --- a/src/purgatory/persistence.rs +++ b/src/purgatory/persistence.rs @@ -106,7 +106,6 @@ pub fn offset_to_instant( #[cfg(test)] mod tests { use super::*; - use std::thread; use std::time::Duration; #[test] @@ -117,17 +116,14 @@ mod tests { let offset = instant_to_offset(future, now_system, now_instant); - // Should be approximately 60 seconds (within tolerance) - assert!(offset.as_secs() >= 59 && offset.as_secs() <= 61); + assert_eq!(offset, Duration::from_secs(60)); } #[test] fn test_instant_to_offset_past() { let now_system = SystemTime::now(); - let past_instant = Instant::now(); - // Simulate some time passing - thread::sleep(Duration::from_millis(10)); let now_instant = Instant::now(); + let past_instant = now_instant - Duration::from_secs(10); let offset = instant_to_offset(past_instant, now_system, now_instant); @@ -137,22 +133,9 @@ mod tests { #[test] fn test_offset_to_instant_with_time_remaining() { - let saved_at = SystemTime::now(); + let saved_at = SystemTime::now() - Duration::from_secs(10); let offset = Duration::from_secs(60); - - // Simulate a very short downtime (< 10ms) - thread::sleep(Duration::from_millis(5)); - - let now_instant = Instant::now(); - let restored = offset_to_instant(offset, saved_at, now_instant); - - // Should be approximately 60 seconds in the future - let remaining = restored.duration_since(now_instant); - assert!( - remaining.as_secs() >= 59 && remaining.as_secs() <= 61, - "Expected ~60s, got {}s", - remaining.as_secs() - ); + assert_restored_within_call_bounds(offset, saved_at, Instant::now()); } #[test] @@ -178,21 +161,23 @@ mod tests { // Convert to offset let offset = instant_to_offset(future, now_system, now_instant); - // Immediately convert back (minimal downtime) - let restored = offset_to_instant(offset, now_system, now_instant); + assert_restored_within_call_bounds(offset, now_system, now_instant); + } - // Should be very close to the original future instant - let diff = if restored > future { - restored.duration_since(future) - } else { - future.duration_since(restored) + // The conversion reads SystemTime internally. Bound that read with actual + // observations instead of assuming a maximum scheduling delay. + fn assert_restored_within_call_bounds( + offset: Duration, + saved_at: SystemTime, + reference: Instant, + ) { + let before = SystemTime::now(); + let restored = offset_to_instant(offset, saved_at, reference); + let after = SystemTime::now(); + let remaining_at = |now: SystemTime| { + offset.saturating_sub(now.duration_since(saved_at).unwrap_or(Duration::ZERO)) }; - - // Allow for small timing differences (< 100ms) - assert!( - diff < Duration::from_millis(100), - "Round trip should preserve instant within 100ms, got {}ms", - diff.as_millis() - ); + assert!(restored >= reference + remaining_at(after)); + assert!(restored <= reference + remaining_at(before)); } } diff --git a/src/purgatory/sync/queue.rs b/src/purgatory/sync/queue.rs index 3226f47..be56f15 100644 --- a/src/purgatory/sync/queue.rs +++ b/src/purgatory/sync/queue.rs @@ -135,17 +135,17 @@ mod tests { entry.attempt_count = 5; entry.next_attempt = Instant::now() + Duration::from_secs(120); - // New event arrives with shorter delay + // Capture the call bounds instead of assuming the scheduler runs promptly. + let original_next = entry.next_attempt; + let before = Instant::now(); entry.on_new_event(Duration::from_secs(10)); + let after = Instant::now(); // Attempt count should be reset assert_eq!(entry.attempt_count, 0); - // next_attempt should be updated to the sooner time - // (within a small tolerance for test timing) - let expected = Instant::now() + Duration::from_secs(10); - assert!(entry.next_attempt <= expected + Duration::from_millis(100)); - assert!(entry.next_attempt >= expected - Duration::from_millis(100)); + assert!(entry.next_attempt >= original_next.min(before + Duration::from_secs(10))); + assert!(entry.next_attempt <= original_next.min(after + Duration::from_secs(10))); } #[test] @@ -160,16 +160,14 @@ mod tests { assert_eq!(entry.attempt_count, 0); // But next_attempt should not be pushed back - assert!(entry.next_attempt <= original_next + Duration::from_millis(100)); + assert_eq!(entry.next_attempt, original_next); } #[test] fn is_ready_checks_both_conditions() { let mut entry = SyncQueueEntry::new(Duration::from_secs(0)); - // Should be ready initially (no delay, not in progress) - // Note: there might be a tiny delay, so we wait a moment - std::thread::sleep(Duration::from_millis(10)); + // A zero delay is already ready; monotonic time cannot move backward. assert!(entry.is_ready()); // Mark as in progress - should not be ready @@ -185,12 +183,13 @@ mod tests { #[test] fn on_sync_complete_increments_and_schedules() { let mut entry = SyncQueueEntry::new(Duration::from_secs(0)); - std::thread::sleep(Duration::from_millis(10)); // Ensure next_attempt has passed entry.in_progress = true; entry.attempt_count = 0; + let before = Instant::now(); entry.on_sync_complete(); + let after = Instant::now(); // Should no longer be in progress assert!(!entry.in_progress); @@ -199,8 +198,7 @@ mod tests { assert_eq!(entry.attempt_count, 1); // Next attempt should be scheduled with backoff (20s for attempt 1) - let expected = Instant::now() + Duration::from_secs(20); - assert!(entry.next_attempt >= expected - Duration::from_millis(100)); - assert!(entry.next_attempt <= expected + Duration::from_millis(100)); + assert!(entry.next_attempt >= before + Duration::from_secs(20)); + assert!(entry.next_attempt <= after + Duration::from_secs(20)); } } diff --git a/src/sync/rejected_index.rs b/src/sync/rejected_index.rs index 718fca3..0808a4e 100644 --- a/src/sync/rejected_index.rs +++ b/src/sync/rejected_index.rs @@ -1507,6 +1507,13 @@ mod tests { keys.sign_event(unsigned).unwrap() } + fn expire_hot_cache(cache: &HotCache) { + let expired_at = Instant::now() - cache.expiry_duration; + for entry in cache.entries.write().unwrap().values_mut() { + entry.cached_at = expired_at; + } + } + fn simulate_checkpoint_downtime(path: &Path, downtime: Duration) { let json = std::fs::read_to_string(path).expect("read rejected-event checkpoint"); let mut state: RejectedCacheState = @@ -1598,7 +1605,7 @@ mod tests { #[tokio::test] async fn test_hot_cache_expires_after_duration() { - let cache = HotCache::new(Duration::from_millis(50)); + let cache = HotCache::new(Duration::from_secs(120)); let event = create_test_event().await; cache.add( @@ -1611,8 +1618,8 @@ mod tests { assert!(cache.contains(&event.id)); - // Wait for expiry - std::thread::sleep(Duration::from_millis(60)); + // Move the cache timestamp to its expiry boundary without waiting. + expire_hot_cache(&cache); let expired = cache.cleanup_expired(); assert_eq!(expired, 1); @@ -1698,14 +1705,21 @@ mod tests { #[tokio::test] async fn test_unrecoverable_ids_expire_with_cold_bound() { - let index = RejectedEventsIndex::new(Duration::from_millis(10), Duration::from_millis(50)); + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); let event = create_test_event().await; index.add_unrecoverable(event.id, 30617); assert!(index.contains(&event.id)); - // Passage of time is the behaviour under test (bounded expiry) - std::thread::sleep(Duration::from_millis(60)); + // Set the rejection timestamp at the cold expiry boundary. + index + .unrecoverable + .entries + .write() + .unwrap() + .get_mut(&event.id) + .unwrap() + .rejected_at = Instant::now() - index.unrecoverable.expiry_duration; assert!(!index.contains(&event.id)); assert_eq!(index.cleanup_expired_unrecoverable(), 1); @@ -1793,8 +1807,8 @@ mod tests { #[tokio::test] async fn test_cleanup_expired_both_tiers() { let index = RejectedEventsIndex::new( - Duration::from_millis(50), // Hot cache expires quickly - Duration::from_millis(100), // Cold index expires slower + Duration::from_secs(120), // Hot cache expires quickly + Duration::from_secs(604800), // Cold index expires slower ); let event = create_test_event().await; @@ -1805,15 +1819,22 @@ mod tests { RejectionReason::DoesNotListService, ); - // Wait for hot cache to expire - std::thread::sleep(Duration::from_millis(60)); + // Expire only the hot tier, independently of scheduler timing. + expire_hot_cache(&index.hot_cache); let (hot_expired, cold_expired) = index.cleanup_expired_for_type("announcement"); assert_eq!(hot_expired, 1); assert_eq!(cold_expired, 0); // Not expired yet - // Wait for cold index to expire - std::thread::sleep(Duration::from_millis(50)); + // Now expire the remaining cold entry explicitly. + index + .cold_index + .entries + .write() + .unwrap() + .get_mut(&event.id) + .unwrap() + .rejected_at = Instant::now() - index.cold_index.expiry_duration; let (hot_expired, cold_expired) = index.cleanup_expired_for_type("announcement"); assert_eq!(hot_expired, 0); // Already cleaned up @@ -1822,8 +1843,7 @@ mod tests { #[tokio::test] async fn test_hot_cache_miss_after_expiry() { - let index = - RejectedEventsIndex::new(Duration::from_millis(50), Duration::from_secs(604800)); + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); let event = create_test_event().await; let pubkey = event.pubkey; let identifier = "test-repo".to_string(); @@ -1835,8 +1855,8 @@ mod tests { RejectionReason::MaintainerNotYetValid, ); - // Wait for hot cache to expire - std::thread::sleep(Duration::from_millis(60)); + // Expire only the hot tier, independently of scheduler timing. + expire_hot_cache(&index.hot_cache); let (removed, hot_events) = index.invalidate_and_get(&pubkey, &identifier, Some(EventType::Announcement)); @@ -1847,8 +1867,7 @@ mod tests { #[tokio::test] async fn test_expired_dependency_candidate_keeps_id_for_targeted_refetch() { - let index = - RejectedEventsIndex::new(Duration::from_millis(50), Duration::from_secs(604800)); + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); let keys = Keys::generate(); let dependency_event = keys .sign_event( @@ -1877,7 +1896,7 @@ mod tests { RejectionReason::Other, ); - std::thread::sleep(Duration::from_millis(60)); + expire_hot_cache(&index.hot_cache); let (event_ids, hot_events) = index.dependency_candidates(&pubkey, &identifier, Some(EventType::State)); @@ -2110,7 +2129,7 @@ mod tests { let state_path = temp_dir.path().join("rejected_cache.json"); let index = RejectedEventsIndex::new( - Duration::from_millis(50), // Hot cache expires quickly + Duration::from_secs(120), // Hot cache expires quickly Duration::from_secs(604800), // Cold index lasts long ); let event = create_test_event().await; @@ -2123,8 +2142,8 @@ mod tests { RejectionReason::MaintainerNotYetValid, ); - // Wait for hot cache to expire - std::thread::sleep(Duration::from_millis(60)); + // Expire only the hot tier, independently of scheduler timing. + expire_hot_cache(&index.hot_cache); index.cleanup_expired_for_type("announcement"); assert_eq!(index.hot_cache_len(), 0); @@ -2135,7 +2154,7 @@ mod tests { // Restore into new index let index2 = - RejectedEventsIndex::new(Duration::from_millis(50), Duration::from_secs(604800)); + RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); index2.restore_from_disk(&state_path).unwrap(); // Verify only cold index restored (hot cache was empty) @@ -2514,8 +2533,8 @@ mod tests { let temp_dir = tempfile::tempdir().unwrap(); let state_path = temp_dir.path().join("rejected_cache.json"); - // Create index with 2 second hot cache expiry - let index = RejectedEventsIndex::new(Duration::from_secs(2), Duration::from_secs(604800)); + // Use the normal TTL and simulate the elapsed part explicitly. + let index = RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); let event = create_test_event().await; index.add_announcement( @@ -2525,24 +2544,40 @@ mod tests { RejectionReason::DoesNotListService, ); - // Wait 200ms (small fraction of TTL) - std::thread::sleep(Duration::from_millis(200)); + // Save an entry that has already consumed part of its TTL. + index + .hot_cache + .entries + .write() + .unwrap() + .get_mut(&event.id) + .unwrap() + .cached_at = Instant::now() - Duration::from_secs(30); // Save to disk index.save_to_disk(&state_path).unwrap(); // Immediately restore (minimal downtime) - let index2 = RejectedEventsIndex::new(Duration::from_secs(2), Duration::from_secs(604800)); + let index2 = + RejectedEventsIndex::new(Duration::from_secs(120), Duration::from_secs(604800)); index2.restore_from_disk(&state_path).unwrap(); - // Event should still be retrievable (has ~1.8s remaining) + // Restoring must retain the elapsed age, not reset the TTL. + assert!( + index2.hot_cache.entries.read().unwrap()[&event.id] + .cached_at + .elapsed() + >= Duration::from_secs(30) + ); + + // The remaining TTL still permits retrieval. let events = index2 .hot_cache .get_maintainer_events(&event.pubkey, "test-repo", None); assert_eq!(events.len(), 1); - // Wait 2 seconds (total 2.2s > 2s expiry) - std::thread::sleep(Duration::from_secs(2)); + // Advance the restored entry to expiry without a wall-clock wait. + expire_hot_cache(&index2.hot_cache); // Now it should be expired let events = index2