From b58691301b20cee9919b8e2e7f8631b4050cdac8 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 21 Sep 2026 15:01:27 +0000 Subject: [PATCH] test(sync): make dependency batching regression deterministic Sequential acceptance publications can cross maintenance ticks and retain independent retry deadlines, so the Darwin test can fail despite correct per-pass batching. Replace that timer-dependent integration assertion with a manager test that seeds both cold dependencies before explicit maintenance. Observe the actual WebSocket REQ and require both IDs in one filter with the correct limit. Bound fixture waits and own its task through an abort guard. Keep existing admission and recovery integration coverage; production sync behavior and timer settings are unchanged. Document the batching precondition. Validation: full workspace suite passed (3,566 passed, 16 ignored), batching regression passed 20 repetitions, and formatting/whitespace checks passed. Independent review found no blockers. Darwin packaging validation is pending. Assisted-by: Codex (GPT-6) --- docs/how-to/test-fixtures.md | 6 + src/sync/mod.rs | 166 ++++++++++++++++++++++++++ tests/sync/maintainer_reprocessing.rs | 106 +--------------- 3 files changed, 173 insertions(+), 105 deletions(-) diff --git a/docs/how-to/test-fixtures.md b/docs/how-to/test-fixtures.md index 138ffe7..9a13144 100644 --- a/docs/how-to/test-fixtures.md +++ b/docs/how-to/test-fixtures.md @@ -72,6 +72,12 @@ and retain assertions about terminal flush and post-push promotion ordering. Fixed sleeps remain appropriate only when elapsed time is itself under test; polling an observable condition must always have a bounded deadline. +For maintenance batching, seed all due inputs before explicitly invoking the +manager pass and inspect the emitted query. Sequential network publications can +straddle timer ticks and acquire different retry deadlines; waiting longer does +not guarantee they will ever share a batch. Keep separate integration coverage +for admission and eventual recovery through the public protocol. + Targeted regressions include `fixture_lifecycle`, `relay_identity`, `git_response_streaming`, and the shared Git-server tests in `sync`. Full package validation remains necessary after these scoped checks. diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 85d2fe5..c7b762a 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -9272,6 +9272,172 @@ impl SyncManager { mod tests { use super::*; + /// Seed both repositories before driving maintenance: sequential network + /// publication can straddle ticks and leave their retry deadlines offset. + #[tokio::test] + async fn unresolved_repositories_share_one_dependency_poll_per_relay() { + use futures_util::{SinkExt, StreamExt}; + use tokio_tungstenite::tungstenite::Message; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("reserve query observer listener"); + let address = listener.local_addr().unwrap(); + let relay_url = format!("ws://{address}"); + let (query_tx, query_rx) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (socket, _) = listener.accept().await.unwrap(); + let mut socket = tokio_tungstenite::accept_async(socket).await.unwrap(); + let mut query_tx = Some(query_tx); + while let Some(frame) = socket.next().await { + let frame = frame.unwrap(); + if !frame.is_text() { + continue; + } + let request: serde_json::Value = + serde_json::from_str(frame.to_text().unwrap()).unwrap(); + if request[0] == "REQ" { + // Observe the complete dependency exchange through its CLOSE. + query_tx + .take() + .expect("one batched dependency query") + .send(request.clone()) + .unwrap(); + socket + .send(Message::Text( + serde_json::json!(["EOSE", request[1]]).to_string().into(), + )) + .await + .unwrap(); + } else if request[0] == "CLOSE" { + break; + } + } + }); + // Abort the owned fixture on assertion failure as well as normal exit. + struct AbortServer(tokio::task::AbortHandle); + impl Drop for AbortServer { + fn drop(&mut self) { + self.0.abort(); + } + } + let _server_guard = AbortServer(server.abort_handle()); + let directory = tempfile::tempdir().expect("create manager directory"); + let git_path = directory.path().join("git"); + let mut config = Config::for_testing(); + config.git_data_path = git_path.to_string_lossy().into_owned(); + config.relay_data_path = directory + .path() + .join("relay") + .to_string_lossy() + .into_owned(); + // Cold-index behavior is the subject; no wall-clock TTL wait is needed. + config.rejected_hot_cache_duration_secs = 0; + let purgatory = Arc::new(crate::purgatory::Purgatory::new(git_path.clone())); + let runtime = crate::nostr::builder::create_relay( + &config, + purgatory.clone(), + crate::grasp06::receive::RepoInitLocks::default(), + None, + ) + .await + .expect("create manager relay runtime"); + // Treat the observer as our own service to suppress unrelated live + // subscriptions. Exact-ID recovery still uses its connected socket. + let mut manager = SyncManager::new( + None, + address.to_string(), + runtime.stores.database.clone(), + runtime.write_policy, + runtime.relay, + &config, + git_path.clone(), + None, + None, + None, + ); + let connection = RelayConnection::new( + relay_url.clone(), + None, + RelayTargetSource::OperatorConfigured, + OutboundTargetPolicy::default(), + ); + connection.connect(5).await.expect("connect query observer"); + manager.connections.insert(relay_url.clone(), connection); + manager.relay_sync_index.write().await.insert( + relay_url.clone(), + RelayState { + connection_status: ConnectionStatus::Connected, + ..RelayState::default() + }, + ); + let invitee = Keys::generate(); + let mut expected_ids = HashSet::new(); + for identifier in ["batched-dependency-one", "batched-dependency-two"] { + let owner = Keys::generate(); + let invitation = EventBuilder::new(Kind::GitRepoAnnouncement, "Owner invitation") + .tags([ + Tag::identifier(identifier), + Tag::custom( + "clone", + [format!("https://example.invalid/{identifier}.git")], + ), + Tag::custom("relays", [relay_url.clone()]), + Tag::custom("maintainers", [invitee.public_key().to_hex()]), + ]) + .finalize(&owner) + .unwrap(); + expected_ids.insert(invitation.id.to_hex()); + manager.rejected_events_index.add_announcement_from_relay( + invitation, + owner.public_key(), + identifier.into(), + RejectionReason::DoesNotListService, + Some(relay_url.clone()), + ); + let acceptance = EventBuilder::new(Kind::GitRepoAnnouncement, "Invitee acceptance") + .tags([ + Tag::identifier(identifier), + Tag::custom("clone", [format!("http://{address}/{identifier}.git")]), + Tag::custom("relays", [relay_url.clone()]), + Tag::custom("maintainers", [owner.public_key().to_hex()]), + ]) + .finalize(&invitee) + .unwrap(); + purgatory.add_announcement( + acceptance, + identifier.into(), + invitee.public_key(), + git_path.join(identifier), + HashSet::from([relay_url.clone()]), + ); + } + + // No background manager loop runs: this pass sees both fresh entries. + tokio::time::timeout( + Duration::from_secs(5), + manager.sync_purgatory_announcements_to_index(), + ) + .await + .expect("maintenance must finish"); + let query = tokio::time::timeout(Duration::from_secs(5), query_rx) + .await + .expect("dependency query deadline") + .expect("query observer"); + assert_eq!(query.as_array().unwrap().len(), 3, "one exact-ID filter"); + let actual_ids: HashSet = serde_json::from_value(query[2]["ids"].clone()).unwrap(); + assert_eq!( + actual_ids, expected_ids, + "both repositories share one query" + ); + assert_eq!(query[2]["limit"], 2); + tokio::time::timeout(Duration::from_secs(5), server) + .await + .expect("query fixture must finish after CLOSE") + .expect("query fixture"); + manager.shutdown().await; + } + #[tokio::test] async fn accepted_dependency_reprocesses_a_synced_policy_orphan() { let directory = tempfile::tempdir().expect("create test directory"); diff --git a/tests/sync/maintainer_reprocessing.rs b/tests/sync/maintainer_reprocessing.rs index d19f116..1d9746b 100644 --- a/tests/sync/maintainer_reprocessing.rs +++ b/tests/sync/maintainer_reprocessing.rs @@ -37,9 +37,7 @@ use std::time::Duration; use nostr_sdk::prelude::*; -use crate::common::{ - censoring_proxy::CensoringProxy, mock_relay::MockRelay, sync_helpers::*, TestRelay, -}; +use crate::common::{sync_helpers::*, TestRelay}; async fn wait_for_log(path: &Path, needle: &str, timeout: Duration) -> bool { let deadline = tokio::time::Instant::now() + timeout; @@ -591,108 +589,6 @@ async fn test_invitee_only_acceptance_recovers_cold_owner_invitation() { owner_relay.stop().await; } -/// Cold dependency polling for separate repositories sharing a relay is -/// combined into one exact-ID query on each maintenance round. -#[tokio::test] -async fn unresolved_repositories_share_one_dependency_poll_per_relay() { - let source = MockRelay::start().await; - let proxy = CensoringProxy::start(source.url()).await; - let invitee_relay = - TestRelay::start_with_sync_and_rejected_hot_cache(Some(proxy.url().to_string()), 1).await; - let invitee_keys = Keys::generate(); - let owners = [Keys::generate(), Keys::generate()]; - let identifiers = ["batched-dependency-one", "batched-dependency-two"]; - - let mut invitations = Vec::new(); - for (owner, identifier) in owners.iter().zip(identifiers) { - let invitation = EventBuilder::new(Kind::GitRepoAnnouncement, "Owner invitation") - .tags([ - Tag::identifier(identifier), - Tag::custom( - "clone", - [format!( - "https://example.invalid/{}/{}.git", - owner.public_key().to_hex(), - identifier - )], - ), - Tag::custom("relays", [proxy.url().to_string()]), - Tag::custom("maintainers", [invitee_keys.public_key().to_hex()]), - ]) - .finalize(owner) - .expect("Failed to create owner invitation"); - send_to_relay_url(source.url(), &invitation) - .await - .expect("Failed to seed owner invitation"); - invitations.push(invitation); - } - - for invitation in &invitations { - let invitation_note = invitation - .id - .to_bech32() - .expect("Failed to encode invitation event ID"); - assert!( - wait_for_log( - &invitee_relay.log_path(), - &invitation_note, - Duration::from_secs(10), - ) - .await, - "Invitee relay should reject and index both owner invitations" - ); - proxy.withhold(invitation.id); - } - - // Passage of the configured one-second hot-cache TTL is required so both - // repositories exercise cold exact-ID recovery on the same relay. - tokio::time::sleep(Duration::from_secs(2)).await; - - let acceptances: Vec = owners - .iter() - .zip(identifiers) - .zip(&invitations) - .map(|((owner, identifier), invitation)| { - repository_announcement( - &invitee_keys, - &[&invitee_relay], - &[owner.public_key()], - identifier, - ) - .custom_created_at(timestamp_after(invitation.created_at)) - .finalize(&invitee_keys) - .expect("Failed to create invitee acceptance") - }) - .collect(); - let client = Client::default(); - client - .add_relay(invitee_relay.url()) - .await - .expect("Failed to add invitee relay"); - client.connect().await; - for acceptance in &acceptances { - client - .send_event(acceptance) - .await - .expect("Failed to publish invitee acceptance"); - } - client.disconnect().await; - - assert!( - wait_for_log( - &invitee_relay.log_path(), - "repository_count=2 requested_count=2 fetched_count=0", - Duration::from_secs(20), - ) - .await, - "Separate repositories should share one dependency query to their common relay" - ); - - invitee_relay.stop().await; - proxy.stop().await; - source.stop().await; -} - /// Listing a maintainer immediately authorizes that maintainer's state for the /// inviting owner's repository; reciprocal acceptance is not required. ///