fix(sync): budget outbound Git commands and speculative fetches

Cold archive sync counted an entire fetch pass as one request, allowing a single missing-tip backlog to issue hundreds of unaccounted commands. Add a shared 60-command sliding domain budget at subprocess admission, covering advertisements, batch fetches, residuals, hedges and integrity repair.

Reserve useful capacity by limiting speculative purgatory requests to two per pass and six per domain per minute, only below half of the command budget. Deferred OIDs remain eligible without entering the miss memo. One-shot integrity repair reserves discovery plus the first fetch atomically, preventing staggered quota expiry from causing advertisement-only retries. It releases unused reservations and retries admission after releasing its storage lease, retaining ordinary priority for accepted data.

Classify explicit Git rate-limit rejections and apply a 60-second domain cooldown; expose admission deferrals in metrics. The accounting unit is a Git command, not an HTTP exchange or a discovered server quota. Keep existing pass admission, purgatory expiry and configured throughput unchanged; adaptive quota tuning and production deployment are outside this change.

Validation: the initial change passed all 994 library tests. Independent review reproduced an integrity retry starvation case; the fix passes all 71 purgatory sync tests, including staggered-window progress and reservation cleanup regressions, and independent re-review is clear. Both fetch integration tests pass, including a fresh read-only archive with 126 missing tips that fetches advertised data and retries deferred OIDs without repeating misses. Workspace/all-target Clippy with warnings denied, formatting and diff checks pass.

Assisted-by: GPT-6
This commit is contained in:
DanConwayDev
2026-10-02 09:38:38 +00:00
parent 9fd6d7c7f0
commit 643367f0b4
11 changed files with 692 additions and 107 deletions
+4
View File
@@ -9,6 +9,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed
- Count individual outbound Git commands against a shared domain budget and
reserve capacity for advertised commits during cold archive sync. Bound
speculative OID fetches and report explicit Git rate-limit rejections.
- Skip paused peers before scheduling purgatory dependency fetches, avoiding
redundant retries and warnings while preserving retry eligibility on recovery.
+8
View File
@@ -497,6 +497,14 @@ pub struct Purgatory {
service
- Foreground Git requests do not draw from this background-only limit
7. **Outbound Git Command Budget**: a shared ledger counts individual
`ls-remote` and fetch commands across purgatory and integrity repair, with
60 commands per domain per minute. Speculative purgatory OIDs get at most
two attempts per pass and six per domain per minute, only below half of
the ordinary command budget. Denied work retries without holding a storage
lease; explicit Git rate-limit rejections trigger a domain cooldown.
See [Git fetch admission](grasp-02-proactive-sync-purgatory-git-data.md#fetch-strategy-advertised-tips-first).
#### Data Types
See [`types.rs`](../../src/purgatory/types.rs) for complete definitions:
@@ -35,7 +35,7 @@ The system scours git servers listed in repository announcements and PR events,
We respect remote server capacity with:
- **Throttling**: Max 5 concurrent requests per domain, 30 requests/minute
- **Throttling**: Max 5 concurrent requests per domain, 60 Git commands/minute
- **Backoff**: Start at 20 seconds, double each attempt, cap at 2 minutes
- **Round-robin**: Fair distribution across repositories waiting for the same domain
- **Fresh start**: New events reset retry count—recent updates often mean fresh data
@@ -45,7 +45,7 @@ We respect remote server capacity with:
### Key Features
✅ **Proactive hunting** - Scours git servers every 2 min (backoff), finds data automatically
✅ **Respectful throttling** - 5 concurrent + 30/min per domain, plays nice with other implementations
✅ **Respectful throttling** - 5 concurrent + 60 commands/min per domain, plays nice with other implementations
✅ **Smart timing** - 3min delay for user pushes, 500ms for synced events
✅ **30min expiry** - Auto-cleanup of events when data never arrives
✅ **Soft expiry for announcements** - Bare repo deleted at 30min, event retained 24h to allow revival
@@ -270,16 +270,16 @@ Git servers have finite resources. Without throttling:
With throttling:
- ✅ Respect server capacity (5 concurrent max per domain)
- ✅ Stay under rate limits (30 requests/min per domain)
- ✅ Bound outbound work (60 Git commands/min per domain)
- ✅ Fair access for all clients
### Two-Level Limits
Each domain has **two independent limits**:
Concurrency and command rate are bounded independently:
#### 1. Concurrent Request Limit (Default: 5)
#### 1. Concurrent Fetch-Pass Limit (Default: 5)
Maximum in-flight requests to a domain at any moment.
Maximum in-flight purgatory fetch passes to a domain at any moment.
**Example**:
@@ -292,23 +292,20 @@ fetch-3 completes → in-flight: 4
Status: HAS CAPACITY (process next queued identifier)
```
#### 2. Rate Limit (Default: 30/min)
#### 2. Rate Limit (Default: 60 commands/min)
Maximum requests in any 60-second sliding window.
Maximum admitted Git commands in any 60-second sliding window, shared
across URLs and including integrity repair. A successful advertised-tip pass
normally spends two units (`ls-remote` and batch fetch). Speculative OIDs have
a smaller allowance within that budget; see [fetch admission](#fetch-strategy-advertised-tips-first).
**Example**:
For example, 30 advertised-tip passes can consume all 60 command units.
Further commands are deferred until earlier starts leave the window. This
quota does not claim to match any particular server's HTTP request limit.
```
t=0s: Request 1 → request_times: [0s]
t=1s: Request 2 → request_times: [0s, 1s]
...
t=30s: Request 30 → request_times: [0s, 1s, ..., 30s]
t=31s: Request 31? → THROTTLED (30 requests in last 60s)
t=61s: Request at t=0s aged out → request_times: [1s, ..., 30s]
t=61s: Request 31 → ALLOWED (only 29 in last 60s)
```
**Implementation**: [`src/purgatory/sync/throttle.rs:DomainThrottle::has_capacity()`](../../src/purgatory/sync/throttle.rs)
**Implementation**: [`git_budget.rs`](../../src/purgatory/sync/git_budget.rs).
The existing [pass scheduler](../../src/purgatory/sync/throttle.rs) also retains
its 60-pass/minute admission ceiling.
### Round-Robin Fairness
@@ -425,8 +422,8 @@ all through the same hardened/pinned git subprocess machinery:
3. **Residual OIDs one at a time** — needed OIDs that were neither
advertised nor discovered as ancestors of the fetched tips are
requested individually. Some servers refuse arbitrary-SHA1 wants, so a
per-OID failure costs exactly one small round trip and never aborts the
rest of the pass.
missing-object rejection does not fail other OIDs. Speculative requests
are bounded as described below.
Missing residual OIDs are memoized per clone URL against a fingerprint of
the sorted advertised OID set. Later passes skip those OIDs while the
@@ -434,9 +431,42 @@ advertisement is unchanged; any ref-tip change clears that URL's memo and
allows them to be retried. Entries expire lazily after 30 minutes and the
memo is capped at 1,024 URLs, evicting the oldest entry when full.
All three phases run beneath the pass's per-domain permit. Separating the
comparison into its own helper in future must retain that permit; `ls-remote`
is outbound work against the same Git server, not an unaccounted preflight.
All three phases run beneath the pass's per-domain permit. That scheduler
retains its five-concurrent-pass and 60-pass/minute limits. A separate shared
command ledger admits **at most 60 Git commands per domain in a sliding
60-second window**: every `ls-remote`, advertised batch, and individual fetch
costs one unit, including failed attempts. Primary, hedge, and integrity-repair
calls share this ledger on the service's `RealSyncContext`. These are Git
commands, not HTTP requests: Git negotiation may perform multiple HTTP
exchanges, so this is a conservative operational budget, not a measured server
quota.
Purgatory residual requests are speculative. They run after advertised tips,
with at most **two per pass**, **six per domain per minute**, and only while
fewer than **30 commands** have been admitted in that domain's current window.
These limits reserve capacity without a new priority queue. They do not
preempt an already-running speculative request or strictly order independent
passes. A cold archive therefore cannot spend an entire command window
crawling hundreds of unadvertised commits from a single event.
Admission denial launches no subprocess and creates no missing-object memo.
Purgatory returns any fetched objects and retries remaining demand through
its existing backoff loop, subject to the ordinary purgatory expiry. It does
not hold a fetch slot or storage lease waiting for the budget. Accepted-data
integrity repairs use ordinary priority even for unadvertised OIDs; their
one-shot caller retries admission after releasing the fetch's family lease.
Integrity retries atomically reserve capacity for discovery and the first
fetch together, preventing staggered expirations from being spent entirely
on repeated advertisements. Starts are charged when commands launch; unused
reservations are released when the pass exits or is cancelled.
An explicit HTTP 429 (or an explicit remote rate-limit diagnostic) logs the
domain, operation, and role and pauses further commands to that domain for
60 seconds. Other failures, including `not our ref`, do not establish a rate
limit. Already-running commands finish normally. This fixed cooldown does not
infer a server quota or parse `Retry-After` headers. Quotas remain unchanged
pending cold-start measurements; increasing them merely to restore the old
unaccounted request volume would also restore its speculative bursts.
### Repository single-flight and delayed hedging
@@ -572,7 +602,7 @@ async fn test_sync_identifier_partial_success() {
.with_fetch_result("https://server1.com/repo.git", Ok(vec!["oid1"]))
.with_fetch_result("https://server2.com/repo.git", Ok(vec!["oid2"]));
let throttle = Arc::new(ThrottleManager::new(5, 30));
let throttle = Arc::new(ThrottleManager::new(5, 60));
let complete = sync_identifier(&mock, "repo", &throttle).await;
assert!(complete); // Both OIDs fetched
@@ -597,7 +627,7 @@ Purgatory sync behavior is configurable via CLI flags or environment variables:
| Setting | CLI Flag | Environment Variable | Default | Description |
| ----------------------- | -------- | -------------------- | ------- | ---------------------------------------------------- |
| Domain concurrent limit | (future) | (future) | `5` | Max concurrent requests per domain |
| Domain rate limit | (future) | (future) | `30` | Max requests per minute per domain |
| Domain rate limit | (future) | (future) | `60` | Max Git commands per minute per domain |
| Sync loop interval | N/A | N/A | `1s` | How often to check for ready identifiers (hardcoded) |
| Default sync delay | N/A | N/A | `180s` | Delay for user-submitted events (hardcoded) |
| Immediate sync delay | N/A | N/A | `500ms` | Delay for sync-triggered events (hardcoded) |
@@ -701,7 +731,13 @@ Failed to fetch OIDs (url=https://server.com/repo.git, error=connection timeout)
Implemented:
- `ngit_purgatory_git_fetch_passes_total` - Outbound git fetch passes
- `ngit_purgatory_git_deferred_total{operation,reason}` - Admission deferrals
(`command_budget`, `speculative_budget`, `pass_budget`, `remote_cooldown`);
counts denial decisions, not the number of unattempted OIDs
- `ngit_purgatory_git_subprocess_total{operation,role,outcome}` - Completed
commands, with explicit rate-limit rejections classified as `rate_limited`
rather than generic failure; domain attribution is in the warning log
- `ngit_purgatory_git_fetch_passes_total` - Completed outbound git fetch passes
(one ls-remote comparison plus fetches per pass)
- `ngit_purgatory_git_fetch_oids_total{kind}` - Per-pass OID outcomes
(`advertised_tip`, `residual_attempted`, `residual_missing`,
@@ -776,7 +812,7 @@ End-to-end tests verify sync behavior with real relay instances:
### 1. Configurable Throttle Limits
**Current**: Hardcoded to 5 concurrent, 30/min per domain
**Current**: Hardcoded to 5 concurrent, 60 commands/min per domain
**Future**: CLI flags `--sync-domain-concurrent` and `--sync-domain-rate-limit`
**Use case**: Operators might want stricter limits for public servers or looser limits for trusted servers.
@@ -820,7 +856,7 @@ The purgatory sync system is a sophisticated, production-ready implementation th
✅ **Batches intelligently** - Groups events by identifier for efficient fetching
✅ **Retries smartly** - Exponential backoff with fresh start on new events
✅ **Throttles respectfully** - 5 concurrent + 30/min per domain, round-robin fairness
✅ **Throttles respectfully** - 5 concurrent + 60 commands/min per domain, round-robin fairness
✅ **Times strategically** - 3min for user events, 500ms for synced events
✅ **Expires responsibly** - 30min auto-cleanup prevents memory leaks
✅ **Soft-expires announcements** - Bare repo deleted at 30min, event retained 24h for revival
+13 -2
View File
@@ -520,14 +520,25 @@ impl FamilyRepairSource for crate::purgatory::sync::RealSyncContext {
url: &str,
oids: &[String],
) -> Result<Vec<String>> {
crate::purgatory::sync::SyncContext::fetch_oids_with_role(
loop {
let result = crate::purgatory::sync::SyncContext::fetch_oids_with_role(
self,
target_view,
url,
oids,
crate::purgatory::sync::GitFetchRole::Integrity,
)
.await
.await;
let Some(deferred) = result.as_ref().err().and_then(|error| {
error.downcast_ref::<crate::purgatory::sync::GitBudgetDeferred>()
}) else {
return result;
};
// Startup/manual repair is a one-shot caller, unlike purgatory's
// retry loop. Wait for admission with the fetch's family lease
// released, then re-check local objects before trying again.
tokio::time::sleep(deferred.retry_after).await;
}
}
}
+15
View File
@@ -261,6 +261,15 @@ lazy_static! {
REGISTRY.register(Box::new(metric.clone())).expect("register metric");
metric
};
static ref PURGATORY_GIT_DEFERRED_TOTAL: CounterVec = {
let metric = CounterVec::new(
Opts::new("ngit_purgatory_git_deferred_total",
"Outbound Git admission deferrals by operation and budget reason"),
&["operation", "reason"],
).expect("build Git admission deferrals counter");
REGISTRY.register(Box::new(metric.clone())).expect("register metric");
metric
};
static ref PURGATORY_GIT_SUBPROCESS_TOTAL: CounterVec = {
let metric = CounterVec::new(
Opts::new(
@@ -358,6 +367,12 @@ pub struct PurgatoryGitSubprocessGuard {
outcome: &'static str,
}
pub fn record_purgatory_git_deferred(operation: &'static str, reason: &'static str) {
PURGATORY_GIT_DEFERRED_TOTAL
.with_label_values(&[operation, reason])
.inc();
}
impl PurgatoryGitSubprocessGuard {
pub fn finish(&mut self, success: bool) {
self.outcome = if success { "success" } else { "failure" };
+174 -19
View File
@@ -240,6 +240,10 @@ use crate::purgatory::Purgatory;
use crate::sync::naughty_list::NaughtyListTracker;
use super::functions::extract_domain;
use super::git_budget::{
is_rate_limit_error, GitBudgetDeferred, GitCommandBudget, GitCommandReservation,
SPECULATIVE_PER_PASS,
};
/// Real implementation of `SyncContext` that connects to actual systems.
///
@@ -274,6 +278,9 @@ pub struct RealSyncContext {
/// advertised ref tips.
miss_memo: Arc<Mutex<HashMap<String, RemoteMissMemo>>>,
/// Counts individual Git commands across every URL, pass, and repair role.
git_budget: GitCommandBudget,
/// Outbound target policy applied before every event-directed git fetch
outbound_policy: OutboundTargetPolicy,
@@ -323,6 +330,7 @@ impl RealSyncContext {
write_policy,
git_naughty_list,
miss_memo: Arc::new(Mutex::new(HashMap::new())),
git_budget: GitCommandBudget::default(),
outbound_policy,
grasp08_peers,
credential_keys,
@@ -676,6 +684,44 @@ async fn drain_git_stream<R: AsyncRead + Unpin>(
Ok(captured)
}
async fn run_budgeted_git_command(
budget: &GitCommandBudget,
reservation: Option<&mut GitCommandReservation<'_>>,
command: Command,
domain: &str,
operation: &'static str,
role: GitFetchRole,
) -> Result<ObservedGitOutput> {
let admission = match reservation {
Some(reservation) => reservation.admit(),
None => budget.admit(
domain,
operation == "fetch_residual" && role != GitFetchRole::Integrity,
),
};
if let Err(reason) = admission {
crate::metrics::record_purgatory_git_deferred(operation, reason);
debug!(domain, operation, reason, "Outbound Git command deferred");
return Err(GitBudgetDeferred {
retry_after: budget.retry_after(domain),
reason,
}
.into());
}
let output = run_observed_git_command(command, domain, operation, role).await?;
if !output.status.success() && is_rate_limit_error(&String::from_utf8_lossy(&output.stderr)) {
budget.rate_limited(domain);
tracing::warn!(
domain,
operation,
role = role.as_str(),
cooldown_secs = 60,
"Git server rate limited outbound requests; deferring this domain"
);
}
Ok(output)
}
async fn run_observed_git_command(
command: Command,
domain: &str,
@@ -846,6 +892,8 @@ async fn run_observed_git_command_with_policy(
crate::metrics::record_purgatory_git_subprocess_output(operation, role, "stderr", stderr_bytes);
if stalled {
metric.finish_with_outcome("stalled");
} else if !status.success() && is_rate_limit_error(&String::from_utf8_lossy(&stderr)) {
metric.finish_with_outcome("rate_limited");
} else {
metric.finish(status.success());
}
@@ -1151,8 +1199,21 @@ impl SyncContext for RealSyncContext {
// advertisement tells us up front which OIDs the remote can
// serve — instead of discovering each missing one through a
// failed "not our ref" upload-pack round trip.
let mut reservation = if role == GitFetchRole::Integrity {
Some(self.git_budget.reserve_pair(&domain).map_err(|reason| {
crate::metrics::record_purgatory_git_deferred("ls_remote", reason);
GitBudgetDeferred {
retry_after: self.git_budget.retry_after(&domain),
reason,
}
})?)
} else {
None
};
let ls_remote_args = vec!["ls-remote".to_string(), url.clone()];
let advertised = match run_observed_git_command(
let advertised = match run_budgeted_git_command(
&self.git_budget,
reservation.as_mut(),
hardened_git_command(
&repo_path,
resolve_pin.as_deref(),
@@ -1189,13 +1250,7 @@ impl SyncContext for RealSyncContext {
stderr
));
}
Err(e) => {
return Err(anyhow::anyhow!(
"git ls-remote command error for {}: {}",
url,
e
))
}
Err(e) => return Err(e.context(format!("git ls-remote command error for {url}"))),
};
let advert_fingerprint = advertised_oids_fingerprint(&advertised);
let memoized_missing = {
@@ -1224,7 +1279,9 @@ impl SyncContext for RealSyncContext {
];
args.extend(advertised_tips.iter().cloned());
match run_observed_git_command(
match run_budgeted_git_command(
&self.git_budget,
reservation.as_mut(),
hardened_git_command(
&repo_path,
resolve_pin.as_deref(),
@@ -1260,13 +1317,7 @@ impl SyncContext for RealSyncContext {
return Err(anyhow::anyhow!("git fetch failed for {}: {}", url, stderr));
}
}
Err(e) => {
return Err(anyhow::anyhow!(
"git fetch command error for {}: {}",
url,
e
))
}
Err(e) => return Err(e.context(format!("git fetch command error for {url}"))),
}
}
@@ -1275,6 +1326,7 @@ impl SyncContext for RealSyncContext {
// time. Arbitrary-SHA1 wants may still be refused by some
// servers, so a per-OID failure must cost exactly one small
// round trip and must not abort the rest.
let mut deferred = None;
let mut residual_attempted = 0usize;
let mut residual_missing = 0usize;
let mut residual_skipped = 0usize;
@@ -1286,7 +1338,10 @@ impl SyncContext for RealSyncContext {
residual_skipped += 1;
continue;
}
residual_attempted += 1;
if role != GitFetchRole::Integrity && residual_attempted >= SPECULATIVE_PER_PASS {
crate::metrics::record_purgatory_git_deferred("fetch_residual", "pass_budget");
break;
}
let args = vec![
"fetch".to_string(),
@@ -1296,7 +1351,9 @@ impl SyncContext for RealSyncContext {
url.clone(),
oid.clone(),
];
match run_observed_git_command(
let result = run_budgeted_git_command(
&self.git_budget,
reservation.as_mut(),
hardened_git_command(
&repo_path,
resolve_pin.as_deref(),
@@ -1309,8 +1366,16 @@ impl SyncContext for RealSyncContext {
"fetch_residual",
role,
)
.await
.await;
if result
.as_ref()
.is_err_and(|error| error.is::<GitBudgetDeferred>())
{
deferred = result.err();
break;
}
residual_attempted += 1;
match result {
Ok(result) if result.status.success() => {}
Ok(result) => {
if result.stderr_truncated {
@@ -1397,6 +1462,11 @@ impl SyncContext for RealSyncContext {
fetched.len(),
);
if role == GitFetchRole::Integrity {
if let Some(error) = deferred {
return Err(error);
}
}
Ok(fetched)
}
@@ -1554,6 +1624,91 @@ mod fetch_helper_tests {
)
}
#[tokio::test]
async fn command_budget_covers_primary_hedge_and_integrity_and_observes_429() {
fn command(script: &str) -> Command {
let mut command = Command::new("sh");
command
.args(["-c", script])
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.kill_on_drop(true);
command
}
let budget = GitCommandBudget::default();
for index in 0..60 {
let role = [
GitFetchRole::Primary,
GitFetchRole::Hedge,
GitFetchRole::Integrity,
][index % 3];
let result = run_budgeted_git_command(
&budget,
None,
command("exit 0"),
"shared",
"ls_remote",
role,
)
.await
.unwrap();
assert!(result.status.success());
}
let denied = run_budgeted_git_command(
&budget,
None,
command("exit 99"),
"shared",
"fetch_batch",
GitFetchRole::Integrity,
)
.await
.err()
.unwrap();
assert_eq!(
denied.downcast_ref::<GitBudgetDeferred>().unwrap().reason,
"command_budget"
);
let rejected = run_budgeted_git_command(
&budget,
None,
command("echo 'fatal: The requested URL returned error: 429' >&2; exit 1"),
"limited",
"fetch_batch",
GitFetchRole::Primary,
)
.await
.unwrap();
assert!(!rejected.status.success());
let cooled = run_budgeted_git_command(
&budget,
None,
command("exit 99"),
"limited",
"ls_remote",
GitFetchRole::Hedge,
)
.await
.err()
.unwrap();
assert_eq!(
cooled.downcast_ref::<GitBudgetDeferred>().unwrap().reason,
"remote_cooldown"
);
assert!(run_budgeted_git_command(
&budget,
None,
command("exit 0"),
"unrelated",
"ls_remote",
GitFetchRole::Primary
)
.await
.unwrap()
.status
.success());
}
#[tokio::test]
async fn stream_drain_bounds_capture_but_counts_all_activity() {
let (mut writer, reader) = tokio::io::duplex(256);
+323
View File
@@ -0,0 +1,323 @@
//! Admission for individual outbound Git commands, shared by purgatory and
//! integrity repair. Denial defers work to the existing retry loop; it never
//! waits while holding a repository's write lease or a fetch-pass slot.
use std::collections::{HashMap, VecDeque};
use std::sync::Mutex;
use std::time::{Duration, Instant};
const WINDOW: Duration = Duration::from_secs(60);
const COMMANDS_PER_MINUTE: usize = 60;
const SPECULATIVE_PER_MINUTE: usize = 6;
pub(super) const SPECULATIVE_PER_PASS: usize = 2;
#[derive(Debug)]
pub(crate) struct GitBudgetDeferred {
pub retry_after: Duration,
pub reason: &'static str,
}
impl std::fmt::Display for GitBudgetDeferred {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "Git command deferred: {}", self.reason)
}
}
impl std::error::Error for GitBudgetDeferred {}
#[derive(Default)]
struct DomainBudget {
starts: VecDeque<(Instant, bool)>,
cooldown_until: Option<Instant>,
reserved: usize,
}
impl DomainBudget {
fn expire(&mut self, now: Instant) {
while self
.starts
.front()
.is_some_and(|(start, _)| now.duration_since(*start) >= WINDOW)
{
self.starts.pop_front();
}
}
}
/// Integrity retries restart discovery. Reserve discovery plus its first fetch
/// together so repeated advertisements cannot consume every newly freed slot.
pub(super) struct GitCommandReservation<'a> {
budget: &'a GitCommandBudget,
domain: &'a str,
remaining: usize,
}
impl GitCommandReservation<'_> {
pub(super) fn admit(&mut self) -> Result<(), &'static str> {
self.admit_at(Instant::now())
}
fn admit_at(&mut self, now: Instant) -> Result<(), &'static str> {
if self.remaining == 0 {
return self.budget.admit_at(self.domain, false, now);
}
let mut domains = self.budget.domains.lock().unwrap();
let budget = domains
.get_mut(self.domain)
.expect("reserved domain remains live");
if budget.cooldown_until.is_some_and(|until| until > now) {
return Err("remote_cooldown");
}
budget.expire(now);
budget.reserved -= 1;
self.remaining -= 1;
budget.starts.push_back((now, false));
Ok(())
}
}
impl Drop for GitCommandReservation<'_> {
fn drop(&mut self) {
if self.remaining > 0 {
self.budget
.domains
.lock()
.unwrap()
.get_mut(self.domain)
.expect("reserved domain remains live")
.reserved -= self.remaining;
}
}
}
#[derive(Default)]
pub(super) struct GitCommandBudget {
domains: Mutex<HashMap<String, DomainBudget>>,
}
impl GitCommandBudget {
pub(super) fn admit(&self, domain: &str, speculative: bool) -> Result<(), &'static str> {
self.admit_at(domain, speculative, Instant::now())
}
pub(super) fn reserve_pair<'a>(
&'a self,
domain: &'a str,
) -> Result<GitCommandReservation<'a>, &'static str> {
self.reserve_pair_at(domain, Instant::now())
}
fn reserve_pair_at<'a>(
&'a self,
domain: &'a str,
now: Instant,
) -> Result<GitCommandReservation<'a>, &'static str> {
let mut domains = self.domains.lock().unwrap();
let budget = domains.entry(domain.to_owned()).or_default();
budget.expire(now);
if budget.cooldown_until.is_some_and(|until| until > now) {
return Err("remote_cooldown");
}
if budget.starts.len() + budget.reserved + 2 > COMMANDS_PER_MINUTE {
return Err("command_budget");
}
budget.reserved += 2;
Ok(GitCommandReservation {
budget: self,
domain,
remaining: 2,
})
}
fn admit_at(&self, domain: &str, speculative: bool, now: Instant) -> Result<(), &'static str> {
let mut domains = self.domains.lock().unwrap();
// Retire inactive domains as clone URLs change over an archive's life.
domains.retain(|_, budget| {
budget
.starts
.back()
.is_some_and(|(last, _)| now.duration_since(*last) < WINDOW)
|| budget.cooldown_until.is_some_and(|until| until > now)
|| budget.reserved > 0
});
let budget = domains.entry(domain.to_owned()).or_default();
budget.expire(now);
if budget.cooldown_until.is_some_and(|until| until > now) {
return Err("remote_cooldown");
}
if budget.starts.len() + budget.reserved >= COMMANDS_PER_MINUTE {
return Err("command_budget");
}
if speculative
&& (budget.starts.len() + budget.reserved >= COMMANDS_PER_MINUTE / 2
|| budget
.starts
.iter()
.filter(|(_, speculative)| *speculative)
.count()
>= SPECULATIVE_PER_MINUTE)
{
return Err("speculative_budget");
}
budget.starts.push_back((now, speculative));
Ok(())
}
pub(super) fn retry_after(&self, domain: &str) -> Duration {
let now = Instant::now();
let domains = self.domains.lock().unwrap();
let Some(budget) = domains.get(domain) else {
return Duration::ZERO;
};
let quota = budget
.starts
.get(1)
.or_else(|| budget.starts.front())
.map(|(start, _)| *start + WINDOW);
quota
.into_iter()
.chain(budget.cooldown_until)
.map(|until| until.saturating_duration_since(now))
.max()
.unwrap_or_default()
}
pub(super) fn rate_limited(&self, domain: &str) {
self.domains
.lock()
.unwrap()
.entry(domain.to_owned())
.or_default()
.cooldown_until = Some(Instant::now() + WINDOW);
}
}
/// Match explicit HTTP rejection diagnostics, not an arbitrary occurrence of
/// "429" in an OID or URL. Other errors do not establish a remote's quota.
pub(super) fn is_rate_limit_error(stderr: &str) -> bool {
stderr.lines().any(|line| {
let lower = line.to_ascii_lowercase();
lower.ends_with("the requested url returned error: 429")
|| lower.trim() == "remote: too many requests"
|| lower.trim() == "remote: rate limit exceeded"
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn staggered_window_reserves_discovery_and_fetch_together() {
let budget = GitCommandBudget::default();
let start = Instant::now();
for second in 0..60 {
assert!(budget
.admit_at("server", false, start + Duration::from_secs(second))
.is_ok());
}
// One newly free slot must not be spent rediscovering refs. Waiting
// for two slots lets the fetch follow discovery even if other work
// competes for admission in between them.
assert!(budget.reserve_pair_at("server", start + WINDOW).is_err());
let now = start + WINDOW + Duration::from_secs(1);
let mut reservation = budget.reserve_pair_at("server", now).unwrap();
assert!(budget.admit_at("server", false, now).is_err());
assert!(reservation.admit_at(now).is_ok());
assert!(budget.admit_at("server", false, now).is_err());
assert!(reservation
.admit_at(now + Duration::from_millis(100))
.is_ok());
assert!(budget
.admit_at("server", false, now + Duration::from_millis(100))
.is_err());
}
#[test]
fn unused_reservations_release_capacity_and_survive_cleanup() {
let budget = GitCommandBudget::default();
let now = Instant::now();
let reservation = budget.reserve_pair_at("server", now).unwrap();
// No command has started yet, but pruning for another domain must
// not remove the outstanding reservation.
assert!(budget.admit_at("other", false, now + WINDOW).is_ok());
drop(reservation);
for _ in 0..60 {
assert!(budget.admit_at("server", false, now + WINDOW).is_ok());
}
assert!(budget.admit_at("server", false, now + WINDOW).is_err());
}
#[test]
fn cold_archive_backlog_reserves_capacity_for_advertised_work() {
let budget = GitCommandBudget::default();
let now = Instant::now();
let mut useful_repositories = 0;
let mut speculative = 0;
// Many repositories share one server. Each has advertised tips and
// hundreds of unavailable tips; the missing backlog cannot spend the
// budget needed by later repositories' advertisements and batches.
for _ in 0..100 {
if budget.admit_at("server", false, now).is_err() {
break;
}
if budget.admit_at("server", false, now).is_err() {
break;
}
useful_repositories += 1;
for _ in 0..126 {
if budget.admit_at("server", true, now).is_err() {
break;
}
speculative += 1;
}
}
assert_eq!(speculative, 6);
assert_eq!(useful_repositories, 27);
assert_eq!(budget.admit_at("server", false, now), Err("command_budget"));
assert!(budget.admit_at("other-server", false, now).is_ok());
assert!(budget.admit_at("server", false, now + WINDOW).is_ok());
assert!(budget.admit_at("server", true, now + WINDOW).is_ok());
}
#[test]
fn speculative_requests_stop_at_half_usage_even_without_prior_speculation() {
let budget = GitCommandBudget::default();
let now = Instant::now();
for _ in 0..30 {
assert!(budget.admit_at("server", false, now).is_ok());
}
assert_eq!(
budget.admit_at("server", true, now),
Err("speculative_budget")
);
assert!(budget.admit_at("server", false, now).is_ok());
}
#[test]
fn explicit_rate_limit_cools_down_all_roles_then_recovers() {
let budget = GitCommandBudget::default();
budget.rate_limited("server");
let now = Instant::now();
for speculative in [false, true] {
assert_eq!(
budget.admit_at("server", speculative, now),
Err("remote_cooldown")
);
assert!(budget.admit_at("other", speculative, now).is_ok());
}
assert!(budget.admit_at("server", false, now + WINDOW).is_ok());
assert!(is_rate_limit_error(
"fatal: unable to access 'https://host/repo': The requested URL returned error: 429\n"
));
assert!(!is_rate_limit_error(
"fatal: https://host/429.git not found"
));
assert!(!is_rate_limit_error(
"fatal: remote error: upload-pack: not our ref 429abc"
));
assert!(!is_rate_limit_error(
"The requested URL returned error: 403"
));
}
}
+2
View File
@@ -9,6 +9,7 @@
mod context;
mod functions;
mod git_budget;
mod r#loop;
mod queue;
mod throttle;
@@ -18,6 +19,7 @@ pub use functions::{
get_throttled_domains_with_untried_urls, sync_identifier, sync_identifier_from_url,
sync_identifier_next_url, ThrottledDomainInfo,
};
pub(crate) use git_budget::GitBudgetDeferred;
pub use queue::SyncQueueEntry;
pub use throttle::{DomainThrottle, ThrottleManager};
+2 -1
View File
@@ -1,4 +1,5 @@
//! Domain-based rate limiting and identifier queue management.
//! Domain-based fetch-pass admission and identifier queue management.
//! Individual command accounting lives in `git_budget`, shared with repair.
//!
//! This module provides per-domain throttling to prevent overwhelming remote
//! git servers during purgatory sync operations. Each domain has:
+3 -2
View File
@@ -480,8 +480,9 @@ impl RelayServer {
));
info!("Git storage and authorization-integrity worker started");
// Create throttle manager for rate limiting remote git servers
// Default: 5 concurrent requests per domain, 60 requests per minute per domain
// Bound scheduled fetch passes (5 concurrent, 60 starts/minute/domain).
// RealSyncContext separately budgets every outbound Git command,
// including commands issued by integrity repair.
let throttle_manager = Arc::new(ThrottleManager::new(5, 60));
throttle_manager.set_context(sync_ctx.clone());
throttle_manager.set_git_naughty_list(git_naughty_list.clone());
+75 -46
View File
@@ -22,7 +22,7 @@
//! The scenario drives the real purgatory sync path end to end: a genuine
//! ngit-grasp relay serves a repository with two real branch tips behind a
//! counting proxy, and the relay under test holds a state event declaring
//! those two tips plus eight tips that exist nowhere (mirroring `market`).
//! those two tips plus missing tips that exist nowhere (mirroring `market`).
//! The proxy records every upload-pack POST with its `want` lines, so the
//! test can assert the shape of the outbound fetching, not just the result.
@@ -41,7 +41,7 @@ use crate::common::{port, MockRelay, TestRelay};
/// Declared ref tips that exist on no reachable server, mirroring the
/// `market` repository's unfetchable state event.
const MISSING_TIP_COUNT: usize = 8;
const MISSING_TIP_COUNT: usize = 2;
/// Wait until every given oid exists in the bare repository at `repo_path`.
async fn wait_for_oids_in_repo(repo_path: &Path, oids: &[&str], deadline: Duration) -> bool {
@@ -66,29 +66,25 @@ async fn wait_for_oids_in_repo(repo_path: &Path, oids: &[&str], deadline: Durati
}
}
/// Wait until the proxy has observed at least `min_total` upload-pack
/// exchanges and the count has been stable for `stable_for`.
async fn wait_for_upload_pack_quiescence(
proxy: &UploadPackCountingProxy,
min_total: usize,
stable_for: Duration,
deadline: Duration,
) -> bool {
let end = tokio::time::Instant::now() + deadline;
let mut last_total = 0usize;
let mut stable_since = tokio::time::Instant::now();
/// Wait for the fetch summary, which is emitted after residual admission and
/// object retention finish. Object arrival alone precedes those operations.
async fn wait_for_completed_passes(relay: &TestRelay, minimum: usize) -> Vec<String> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
loop {
let total = proxy.exchanges().len();
if total != last_total {
last_total = total;
stable_since = tokio::time::Instant::now();
} else if total >= min_total && stable_since.elapsed() >= stable_for {
return true;
let log = std::fs::read_to_string(relay.log_path()).unwrap_or_default();
let passes: Vec<String> = log
.lines()
.filter(|line| line.contains("Purgatory git fetch pass complete"))
.map(str::to_owned)
.collect();
if passes.len() >= minimum {
return passes;
}
if tokio::time::Instant::now() >= end {
return false;
}
tokio::time::sleep(Duration::from_millis(200)).await;
assert!(
tokio::time::Instant::now() < deadline,
"fetch pass should complete"
);
tokio::time::sleep(Duration::from_millis(50)).await;
}
}
@@ -99,7 +95,7 @@ async fn wait_for_upload_pack_quiescence(
/// 2. A counting proxy fronts the source's git smart-HTTP endpoint.
/// 3. The relay under test bootstrap-syncs from a MockRelay serving a
/// later announcement whose clone tag points at the proxy, plus a state
/// event declaring the two real tips and eight tips that exist nowhere.
/// event declaring the two real tips and missing tips that exist nowhere.
/// Both events sit in purgatory; the purgatory sync loop fetches
/// through the proxy.
/// 4. The two real tips must arrive, and the outbound request shape must
@@ -110,6 +106,15 @@ async fn wait_for_upload_pack_quiescence(
/// single missing object must never fail a batch).
#[tokio::test]
async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
exercise_fetch_backlog(MISSING_TIP_COUNT, false).await;
}
#[tokio::test]
async fn cold_archive_fetches_advertised_tips_without_crawling_missing_backlog() {
exercise_fetch_backlog(126, true).await;
}
async fn exercise_fetch_backlog(missing_tip_count: usize, archive: bool) {
// 1. Source relay + counting proxy in front of its git endpoint, and
// the MockRelay that will carry the events for the relay under test.
let source = TestRelay::start().await;
@@ -208,10 +213,10 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
.finalize(&keys)
.expect("sign syncing announcement");
let missing_tips: Vec<String> = (0..MISSING_TIP_COUNT)
let missing_tips: Vec<String> = (0..missing_tip_count)
.map(|index| format!("beef{index:036x}"))
.collect();
let missing_branch_names: Vec<String> = (0..MISSING_TIP_COUNT)
let missing_branch_names: Vec<String> = (0..missing_tip_count)
.map(|index| format!("missing-{index}"))
.collect();
let mut branches: Vec<(&str, &str)> =
@@ -247,10 +252,12 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
.expect("send state event to mock relay");
// Negentropy is disabled because MockRelay does not support NIP-77.
let syncing = TestRelay::start_on_reservation_with_options(
let syncing = TestRelay::start_on_reservation_with_archive_and_sync(
syncing_reservation,
Some(mock.url().to_string()),
true,
archive,
archive,
)
.await;
@@ -274,8 +281,7 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
// 6. Let the sync pass finish so requests aimed at the missing tips
// (however the client shapes them) are all recorded.
wait_for_upload_pack_quiescence(&proxy, 1, Duration::from_secs(3), Duration::from_secs(60))
.await;
let first_passes = wait_for_completed_passes(&syncing, 1).await;
// 7. Regression assertions on the outbound request shape.
let exchanges = proxy.exchanges();
@@ -326,7 +332,8 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
// 9. A fresh replaceable state event resets the queue backoff and
// deterministically triggers another pass. With an unchanged remote
// advertisement, the miss memo must suppress every failed fetch.
// advertisement, the miss memo must suppress attempted misses while
// deferred OIDs remain eligible.
let failed_after_first_pass = proxy
.exchanges()
.iter()
@@ -343,32 +350,54 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
.await
.expect("send fresh state event to mock relay");
let second_pass_deadline = tokio::time::Instant::now() + Duration::from_secs(60);
while proxy.info_refs_count() <= info_refs_before_second_pass {
let passes = wait_for_completed_passes(&syncing, first_passes.len() + 1).await;
assert!(proxy.info_refs_count() > info_refs_before_second_pass);
for pass in &passes {
let attempted = pass
.split("residual_attempted=")
.nth(1)
.unwrap()
.split_whitespace()
.next()
.unwrap()
.parse::<usize>()
.unwrap();
assert!(
tokio::time::Instant::now() < second_pass_deadline,
"fresh state event should trigger a second ls-remote"
attempted <= 2,
"each pass must bound speculative work: {pass}"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
assert!(
wait_for_upload_pack_quiescence(
&proxy,
proxy.exchanges().len(),
Duration::from_secs(3),
Duration::from_secs(60),
)
.await,
"second fetch pass should become quiescent"
);
let failed_after_second_pass = proxy
.exchanges()
.iter()
.filter(|exchange| !exchange.ok)
.count();
if archive {
assert!(
failed_after_second_pass > failed_after_first_pass,
"deferred OIDs must remain eligible on the next pass"
);
assert!(
failed_after_second_pass <= 6,
"cold backlog must respect the shared speculative allowance"
);
} else {
assert_eq!(
failed_after_second_pass, failed_after_first_pass,
"unchanged advertisement must produce no new failed upload-pack exchanges"
"unchanged advertisement must suppress already attempted missing OIDs"
);
}
let failed_oids: Vec<String> = proxy
.exchanges()
.iter()
.filter(|exchange| exchange.not_our_ref)
.flat_map(|exchange| exchange.wants.clone())
.collect();
let unique: std::collections::HashSet<_> = failed_oids.iter().collect();
assert_eq!(
unique.len(),
failed_oids.len(),
"memoized misses must not be requested again"
);
source_client.disconnect().await;