mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
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)
This commit is contained in:
@@ -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.
|
||||
|
||||
+166
@@ -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<String> = 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");
|
||||
|
||||
@@ -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<Event> = 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.
|
||||
///
|
||||
|
||||
Reference in New Issue
Block a user