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. ///