mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
fix(sync): keep relay retries scheduler-owned
The bounded connection workers only provide a real global limit if they own every retry. nostr-sdk enables an independent reconnect loop by default, so a relay that disconnected after its first successful invitation sync could reconnect outside the semaphore, health backoff, and actor metrics. Detached workers could also outlive an aborted sync manager while waiting for capacity or a socket. Disable SDK auto-reconnect for managed relay connections and have every queued or active worker observe the manager's shutdown channel. If its result receiver disappears, a worker closes the connection instead of leaving an unowned socket behind.
This commit is contained in:
@@ -29,6 +29,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
||||
- Fix large invitation relay sets repeatedly rebuilding their subscriptions by
|
||||
applying the 70-filter consolidation threshold only to fragmentation above
|
||||
the relay's irreducible desired live-filter baseline.
|
||||
- Keep invitation relay connection ownership inside the bounded scheduler by
|
||||
disabling the SDK's independent auto-reconnect loop and cancelling queued or
|
||||
active connection workers when the sync manager shuts down.
|
||||
- Sideband-aware Git clients now receive periodic progress while GRASP performs post-push purgatory promotion and cross-owner repository alignment, preventing the client I/O timeout from expiring during unusually complex finalization.
|
||||
- Smart HTTP pushes now expose the terminal receive-pack flush only after GRASP has finished promoting the matching repository announcement and state from purgatory. Git progress remains streamed while large packs are resolved and checked, but an immediately following clone or proposal push can now rely on a completed push being queryable on the relay.
|
||||
- Batch the one-time deletion-request lifecycle migration so large production databases do not remain unavailable while LMDB commits every historical request in separate transactions.
|
||||
|
||||
@@ -27,7 +27,8 @@ Key Architectural Points:
|
||||
- PendingBatch tracks each new set of filters that may require pagination until they are complete
|
||||
- Websocket handshakes run in at most eight bounded workers outside the sync
|
||||
actor; only the actor applies their results, and subscriptions start only
|
||||
after the relay reports `Connected`
|
||||
after the relay reports `Connected`. The sync manager owns retry/backoff
|
||||
rather than the SDK, and shutdown cancels queued or active workers
|
||||
- Recompute desired filters when connection is established to ensure filters are as consolidated as possible
|
||||
- Consolidation bounds incremental fragmentation to 70 subscriptions above
|
||||
the irreducible desired live-filter baseline, so large desired sets remain
|
||||
|
||||
+29
-10
@@ -2580,20 +2580,39 @@ impl SyncManager {
|
||||
let health_tracker = Arc::clone(&self.health_tracker);
|
||||
let semaphore = Arc::clone(&self.connect_attempt_semaphore);
|
||||
let timeout = self.health_tracker.base_backoff_secs();
|
||||
let Some(mut shutdown_rx) = self.shutdown_tx.as_ref().map(|sender| sender.subscribe())
|
||||
else {
|
||||
tracing::error!(relay = %relay_url, "Connection scheduler has no shutdown signal");
|
||||
return;
|
||||
};
|
||||
tokio::spawn(async move {
|
||||
let Some(_permit) = begin_connect_attempt(semaphore, health_tracker, &relay_url).await
|
||||
else {
|
||||
let permit = tokio::select! {
|
||||
permit = begin_connect_attempt(semaphore, health_tracker, &relay_url) => permit,
|
||||
_ = shutdown_rx.recv() => return,
|
||||
};
|
||||
let Some(_permit) = permit else {
|
||||
return;
|
||||
};
|
||||
let outcome = match connection.connect(timeout).await {
|
||||
Ok(()) => ConnectAttemptOutcome::Connected,
|
||||
Err(error) => ConnectAttemptOutcome::Failed(error),
|
||||
let outcome = tokio::select! {
|
||||
result = connection.connect(timeout) => match result {
|
||||
Ok(()) => ConnectAttemptOutcome::Connected,
|
||||
Err(error) => ConnectAttemptOutcome::Failed(error),
|
||||
},
|
||||
_ = shutdown_rx.recv() => {
|
||||
connection.disconnect().await;
|
||||
return;
|
||||
}
|
||||
};
|
||||
let _ = result_tx.send(ConnectAttemptResult {
|
||||
relay_url,
|
||||
token,
|
||||
outcome,
|
||||
});
|
||||
if result_tx
|
||||
.send(ConnectAttemptResult {
|
||||
relay_url,
|
||||
token,
|
||||
outcome,
|
||||
})
|
||||
.is_err()
|
||||
{
|
||||
connection.disconnect().await;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -208,6 +208,7 @@ impl RelayConnection {
|
||||
let relay = tokio::time::timeout_at(connection_deadline, async {
|
||||
self.client
|
||||
.add_relay(&self.url)
|
||||
.reconnect(false)
|
||||
.await
|
||||
.map_err(|e| format!("Failed to add relay {}: {}", self.url, e))?;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user