mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
The SelfSubscriber feeds the sync manager's repository index from the service's own accepted events. It dialled our own public WebSocket endpoint with an unauthenticated client, which a private instance's NIP-42 gate correctly refused: the gate runs in the HTTP layer, ahead of LocalRelay, and cannot distinguish our own process from any other anonymous dialler. Announcements accepted at runtime therefore never reached the sync index until a restart rebuilt it from the database, stalling proactive sync and the dynamic membership derived from accepted relay owners — worst on exactly the private-to-private mirroring GRASP-08 exists to enable. Attach the subscriber to the embedded relay in-process instead. A custom WebSocketTransport hands LocalRelay one end of an in-memory duplex pair and keeps the other, so the client gets an ordinary relay session with the same framing, subscription handling, and post-save broadcast, minus the listener, the auth gate, and the network round trip. This replaces the loopback dial in both modes: the public-mode feed is identical in content and strictly more reliable, and it removes a self-directed reconnect loop. Chosen over attaching the relay owner key as a NIP-42 authenticator plus adding the owner pubkey to the effective member set. That alternative works, but widens the member set and the authenticated surface to solve a problem that is not authentication: there is no remote party here. The in-process route needs no key, no membership entry, and no configuration, and nothing reaches the subscriber that the relay did not already accept and persist. Correctness assumptions: LocalRelay applies no NIP-42 or query policy of its own — private-mode access control lives entirely in the HTTP layer — so an in-process session is exactly a local session, not a bypassed remote one. Both ends speak raw WebSocket framing over the duplex with no HTTP upgrade, matching take_connection's Role::Server. The session consumes one connection permit, as the loopback dial did. The GRASP-08 regression test no longer restarts the relay over persistent LMDB: it publishes an announcement at runtime and asserts sync connections to both referenced relays, which is only possible if the live feed reached the index. Verified to fail against the previous implementation (60s deadline, no connections) and pass with this one. req_concurrency's source relay now applies the production outbound target policy. Its scenario lists a proxy URL in the announcement, and with a reliable live feed the source discovers that URL — a distinct host:port that happens to front itself — as an event-directed sync target and opens its own REQ traffic through the proxy, contending for a budget the test means to measure for the syncing relay alone. That behavior is pre-existing and was already reachable after a restart; only its timing changed. The policy keeps the source scenery without weakening the assertion. Deliberately excluded: neg_concurrency shares that topology but passes unchanged, so its fixture is left alone; the duplicate NIP-11 fetch between the pre-dial preflight probe and the post-connect hint fetch is untouched. Validation: cargo clippy --all-targets -D warnings; cargo test --lib (793 passed); cargo test --test private_mode --test sync --test outbound_policy --test purgatory_sync (255 passed, including the full sync suite under parallel load). req_concurrency's startup-burst test passed 4/4 isolated runs after the fixture change.
293 lines
12 KiB
Rust
293 lines
12 KiB
Rust
//! Transient REQ Concurrency Bound Tests
|
||
//!
|
||
//! Regression coverage for a production failure observed on gitnostr.com:
|
||
//! historic REQ+EOSE sync subscribed every byte-budgeted filter group of a
|
||
//! batch in one fast loop and left them all open concurrently until EOSE.
|
||
//! On strfry-family relays REQ subscriptions share `maxSubsPerConnection`
|
||
//! with negentropy views and live subscriptions; nos.lol (budget 20)
|
||
//! answered each gitnostr.com startup with a burst of
|
||
//! "ERROR: too many concurrent REQs" NOTICEs (6 in one second at the
|
||
//! 2026-08-05 startup; 7,739 over the prior three days).
|
||
//!
|
||
//! The scenario drives the real sync path end to end: a genuine ngit-grasp
|
||
//! source relay sits behind a proxy that mimics the strfry limit, rejecting
|
||
//! REQ frames beyond 5 concurrent transient subscriptions. The syncing
|
||
//! relay runs with negentropy disabled so historic sync takes the REQ+EOSE
|
||
//! path, and the watched set is sized so that one startup batch needs more
|
||
//! REQ groups than the limit (2500 root events → six byte-budgeted chunks →
|
||
//! 18 e/E/q filters; two full 32 KB filters fill a 96 KB REQ group, giving
|
||
//! ~8 groups).
|
||
//!
|
||
//! The burst only occurs when startup recomputes filters from an
|
||
//! already-populated index — during initial discovery, root events trickle
|
||
//! in through short batching windows and batches stay small. The scenario
|
||
//! therefore syncs everything into persistent storage first, then restarts
|
||
//! the relay, exactly the production startup shape. Unbounded subscription
|
||
//! creation draws rejections at that startup; the bounded client must draw
|
||
//! none while still overlapping REQs and completing the sync.
|
||
//!
|
||
//! See docs/explanation/sync-scaling-constraints.md for the budget model.
|
||
|
||
use std::time::Duration;
|
||
|
||
use nostr_sdk::prelude::*;
|
||
|
||
use crate::common::purgatory_helpers::{
|
||
create_state_event, create_test_repo_with_commit, push_to_relay, CommitVariant,
|
||
};
|
||
use crate::common::req_limiting_proxy::ReqLimitingProxy;
|
||
use crate::common::sync_helpers::{repo_coord, send_to_relay, wait_for_event_on_relay, TestClient};
|
||
use crate::common::{port, TestRelay};
|
||
|
||
/// Concurrency limit enforced by the proxy.
|
||
///
|
||
/// Matches `MAX_CONCURRENT_TRANSIENT_REQS` in `sync::relay_connection`: the
|
||
/// test relay must stay within its own advertised bound, so a relay
|
||
/// enforcing exactly that bound must never reject a REQ.
|
||
const PROXY_REQ_LIMIT: usize = 5;
|
||
|
||
/// Root events to seed. Byte-budgeted chunking fits ~489 hex event IDs per
|
||
/// filter (32 KB of serialized tag values), so 2500 IDs split into six
|
||
/// chunks; six chunks × three tag variants (e/E/q) give 18 filters, and two
|
||
/// full 32 KB filters fill a 96 KB REQ group — roughly eight groups in a
|
||
/// single historic batch, more concurrent REQ subscriptions than the bound,
|
||
/// so bounded scheduling is required to avoid rejections.
|
||
const ISSUE_COUNT: usize = 2500;
|
||
|
||
/// Wait until the proxy has seen at least `min_total` transient REQs
|
||
/// (opened + rejected) and the count has been stable for `stable_for`.
|
||
///
|
||
/// Returns the final (opened, rejected) pair, or `None` on deadline.
|
||
async fn wait_for_req_quiescence(
|
||
proxy: &ReqLimitingProxy,
|
||
min_total: usize,
|
||
stable_for: Duration,
|
||
deadline: Duration,
|
||
) -> Option<(usize, usize)> {
|
||
let end = tokio::time::Instant::now() + deadline;
|
||
let mut last_total = 0usize;
|
||
let mut stable_since = tokio::time::Instant::now();
|
||
loop {
|
||
let opened = proxy.opened_count();
|
||
let rejected = proxy.rejected_count();
|
||
let total = opened + rejected;
|
||
if total != last_total {
|
||
last_total = total;
|
||
stable_since = tokio::time::Instant::now();
|
||
} else if total >= min_total && stable_since.elapsed() >= stable_for {
|
||
return Some((opened, rejected));
|
||
}
|
||
if tokio::time::Instant::now() >= end {
|
||
return None;
|
||
}
|
||
tokio::time::sleep(Duration::from_millis(200)).await;
|
||
}
|
||
}
|
||
|
||
/// Scenario:
|
||
/// 1. Source relay hosts one repository (announcement + state + git data)
|
||
/// and 2500 historic issues, each a root event of that repository.
|
||
/// 2. A proxy in front of the source enforces the strfry-style limit of 5
|
||
/// concurrent transient REQs and records peak concurrency and
|
||
/// rejections.
|
||
/// 3. The syncing relay (negentropy disabled, persistent storage)
|
||
/// bootstraps through the proxy and syncs all issues.
|
||
/// 4. The syncing relay restarts. Startup recomputes sync filters from the
|
||
/// persisted index of 2500 root events, producing one historic batch
|
||
/// with ~8 REQ groups — more than the limit.
|
||
/// 5. The relay must complete that startup sync without a single rejection
|
||
/// while still keeping REQ subscriptions overlapping.
|
||
#[tokio::test]
|
||
async fn startup_historic_sync_stays_within_relay_req_concurrency_limit() {
|
||
// 1. Pre-allocate the syncing relay port for announcement tags.
|
||
let syncing_reservation = port::reserve_port();
|
||
let syncing_domain = format!("127.0.0.1:{}", syncing_reservation.port());
|
||
|
||
// 2. Source relay with the strfry-style limiting proxy in front. The
|
||
// announcement's relays tag lists the proxy (so the syncing relay
|
||
// targets the proxied connection), not the source itself, so the
|
||
// source runs in archive-all mode to accept it.
|
||
// The source applies the production outbound target policy so it does
|
||
// not itself sync through the proxy that fronts it: the proxy's REQ
|
||
// budget belongs to the syncing relay under test.
|
||
let source = TestRelay::start_archive_source_behind_proxy().await;
|
||
let proxy = ReqLimitingProxy::start(source.url(), PROXY_REQ_LIMIT).await;
|
||
|
||
// 3. One hosted repository: announcement + state event + git data. The
|
||
// relays tag lists the proxy so the syncing relay targets the proxied
|
||
// connection for this repository's sync work.
|
||
let git_temp_dir = tempfile::tempdir().expect("create temp dir for git repo");
|
||
let commit_hash = create_test_repo_with_commit(git_temp_dir.path(), CommitVariant::StateTest)
|
||
.expect("create test git repo");
|
||
|
||
let keys = Keys::generate();
|
||
let npub = keys.public_key().to_bech32().expect("npub");
|
||
let clone_urls = vec![
|
||
format!("http://{}/{}/req-repo.git", source.domain(), npub),
|
||
format!("http://{}/{}/req-repo.git", syncing_domain, npub),
|
||
];
|
||
let relay_urls = vec![proxy.url().to_string(), format!("ws://{}", syncing_domain)];
|
||
|
||
let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "Repository state")
|
||
.tags(vec![
|
||
Tag::identifier("req-repo"),
|
||
Tag::custom("clone", clone_urls.clone()),
|
||
Tag::custom("relays", relay_urls.clone()),
|
||
])
|
||
.finalize(&keys)
|
||
.expect("sign repo announcement");
|
||
let state_event = create_state_event(
|
||
&keys,
|
||
"req-repo",
|
||
&[("main", &commit_hash)],
|
||
&[],
|
||
&clone_urls.iter().map(|s| s.as_str()).collect::<Vec<_>>(),
|
||
&relay_urls.iter().map(|s| s.as_str()).collect::<Vec<_>>(),
|
||
)
|
||
.expect("create state event");
|
||
|
||
send_to_relay(&source, &announcement)
|
||
.await
|
||
.expect("send announcement to source");
|
||
send_to_relay(&source, &state_event)
|
||
.await
|
||
.expect("send state event to source");
|
||
push_to_relay(git_temp_dir.path(), &source.domain(), &npub, "req-repo")
|
||
.expect("push git data to source relay");
|
||
|
||
// The push releases the announcement from purgatory; wait until the
|
||
// source actually serves it before seeding dependent events.
|
||
assert!(
|
||
wait_for_event_on_relay(
|
||
source.url(),
|
||
Filter::new().id(announcement.id),
|
||
Duration::from_secs(15),
|
||
)
|
||
.await,
|
||
"announcement should be released from purgatory on the source"
|
||
);
|
||
|
||
// 4. Seed historic issues; each becomes a tracked root event once
|
||
// synced. One persistent client keeps seeding fast, and distinct
|
||
// descending created_at timestamps let REQ+EOSE pagination cursors
|
||
// (until = oldest seen) terminate instead of re-fetching the same
|
||
// page forever.
|
||
let coordinate = repo_coord(&keys, "req-repo");
|
||
let base_created_at = Timestamp::now().as_secs() - ISSUE_COUNT as u64 - 10;
|
||
|
||
// The embedded relay rate-limits notes per minute PER CONNECTION with a
|
||
// full initial token bucket, so seed through a fresh connection per
|
||
// batch, each batch staying under the initial burst allowance.
|
||
const SEED_BATCH: usize = 50;
|
||
let mut issue_ids = Vec::with_capacity(ISSUE_COUNT);
|
||
for batch_start in (0..ISSUE_COUNT).step_by(SEED_BATCH) {
|
||
let seed_client = TestClient::new(source.url(), Keys::generate())
|
||
.await
|
||
.expect("connect seed client to source");
|
||
for index in batch_start..(batch_start + SEED_BATCH).min(ISSUE_COUNT) {
|
||
let issue = EventBuilder::new(Kind::GitIssue, format!("Historic issue {index}"))
|
||
.tags(vec![Tag::custom("a", vec![coordinate.clone()])])
|
||
.custom_created_at(Timestamp::from(base_created_at + index as u64))
|
||
.finalize(&keys)
|
||
.expect("sign issue event");
|
||
issue_ids.push(issue.id);
|
||
seed_client
|
||
.send_event(&issue)
|
||
.await
|
||
.expect("send issue to source");
|
||
}
|
||
seed_client.disconnect().await;
|
||
}
|
||
|
||
// Seeding sanity: the last issue must be queryable on the source before
|
||
// the syncing relay is pointed at it.
|
||
assert!(
|
||
wait_for_event_on_relay(
|
||
source.url(),
|
||
Filter::new().id(issue_ids[ISSUE_COUNT - 1]),
|
||
Duration::from_secs(15),
|
||
)
|
||
.await,
|
||
"seeded issues should be accepted by the source relay"
|
||
);
|
||
|
||
// 5. Phase 1: the syncing relay (negentropy disabled) bootstraps
|
||
// through the proxy and syncs everything into persistent storage.
|
||
// Root events trickle in through short batching windows here, so
|
||
// historic batches stay small; the burst this test guards against
|
||
// comes at the next startup.
|
||
let syncing_git_dir = tempfile::tempdir().expect("create syncing git dir");
|
||
let syncing_data_dir = tempfile::tempdir().expect("create syncing relay data dir");
|
||
let syncing = TestRelay::start_on_reservation_persistent_sync(
|
||
syncing_reservation,
|
||
Some(proxy.url().to_string()),
|
||
true,
|
||
syncing_git_dir.path().to_path_buf(),
|
||
syncing_data_dir.path().to_path_buf(),
|
||
)
|
||
.await;
|
||
|
||
// 6. The issues must arrive via historic sync (sampled ends of the
|
||
// range cover the byte-budgeted chunks).
|
||
for issue_id in [
|
||
issue_ids[0],
|
||
issue_ids[ISSUE_COUNT / 2],
|
||
issue_ids[ISSUE_COUNT - 1],
|
||
] {
|
||
assert!(
|
||
wait_for_event_on_relay(
|
||
syncing.url(),
|
||
Filter::new().id(issue_id),
|
||
Duration::from_secs(120),
|
||
)
|
||
.await,
|
||
"issue should reach the syncing relay through historic sync"
|
||
);
|
||
}
|
||
|
||
// 7. Let phase-1 REQ activity settle, then restart. Startup recomputes
|
||
// sync filters from the full persisted index (2500 root events),
|
||
// producing one historic batch with ~8 byte-budgeted REQ groups —
|
||
// the production startup burst.
|
||
wait_for_req_quiescence(&proxy, 1, Duration::from_secs(3), Duration::from_secs(120))
|
||
.await
|
||
.expect("phase-1 REQ activity should settle before restart");
|
||
let before_restart = proxy.opened_count() + proxy.rejected_count();
|
||
|
||
let syncing = syncing.restart().await;
|
||
|
||
// 8. Wait for the startup burst to complete and settle: the restarted
|
||
// relay re-syncs the repository batch and the ~8 root-event groups.
|
||
let (opened, rejected) = wait_for_req_quiescence(
|
||
&proxy,
|
||
before_restart + 7,
|
||
Duration::from_secs(3),
|
||
Duration::from_secs(120),
|
||
)
|
||
.await
|
||
.expect("restart REQ burst should complete and settle");
|
||
|
||
// 9. The regression assertions (cumulative across both phases; the
|
||
// trickle-shaped phase 1 must stay within the bound too).
|
||
assert_eq!(
|
||
rejected,
|
||
0,
|
||
"relay exceeded the per-connection REQ concurrency limit \
|
||
(opened: {opened}, peak: {})",
|
||
proxy.peak_concurrent()
|
||
);
|
||
assert!(
|
||
proxy.peak_concurrent() <= PROXY_REQ_LIMIT,
|
||
"peak concurrent transient REQs {} exceeded the bound {PROXY_REQ_LIMIT}",
|
||
proxy.peak_concurrent()
|
||
);
|
||
assert!(
|
||
opened > PROXY_REQ_LIMIT,
|
||
"scenario must open more REQs than the bound to exercise queueing, opened: {opened}"
|
||
);
|
||
|
||
syncing.stop().await;
|
||
proxy.stop().await;
|
||
source.stop().await;
|
||
}
|