mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
Merge #03609f62: Budget outbound Git commands and speculative fetches
nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqsqxcylvt0rvp9fg2le2wctnrs9p0eycpzte93utway6qy7vcp354gmcqaw0 PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 PR description: A cold archive could issue hundreds of speculative Git fetches while the domain limiter counted the entire pass as one request. Count each outbound Git command against a shared 60/minute/domain budget, and limit speculative purgatory fetches to two per pass and six per minute while less than half the budget is used. Deferred commits remain eligible for later passes. Primary, hedge and integrity requests share accounting. Integrity reserves discovery plus its first fetch together so staggered quota expiry cannot trap repair in repeated advertisements. Explicit rate-limit rejections trigger a 60-second domain cooldown and a distinct metric outcome; admission deferrals have their own counter. The budget counts Git commands, not individual HTTP exchanges. Independent review found and reproduced the integrity starvation case; it was fixed and the re-review is clear. Validation: 71 purgatory sync unit tests and both fetch integration tests pass, including a fresh read-only archive with 126 unavailable tips. The original change passed all 994 library tests. Workspace/all-target Clippy with warnings denied and formatting are checked. Production cold-start quota tuning remains follow-up.
This commit is contained in:
@@ -9,6 +9,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
### Fixed
|
### 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
|
- Skip paused peers before scheduling purgatory dependency fetches, avoiding
|
||||||
redundant retries and warnings while preserving retry eligibility on recovery.
|
redundant retries and warnings while preserving retry eligibility on recovery.
|
||||||
|
|
||||||
|
|||||||
@@ -497,6 +497,14 @@ pub struct Purgatory {
|
|||||||
service
|
service
|
||||||
- Foreground Git requests do not draw from this background-only limit
|
- 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
|
#### Data Types
|
||||||
|
|
||||||
See [`types.rs`](../../src/purgatory/types.rs) for complete definitions:
|
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:
|
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
|
- **Backoff**: Start at 20 seconds, double each attempt, cap at 2 minutes
|
||||||
- **Round-robin**: Fair distribution across repositories waiting for the same domain
|
- **Round-robin**: Fair distribution across repositories waiting for the same domain
|
||||||
- **Fresh start**: New events reset retry count—recent updates often mean fresh data
|
- **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
|
### Key Features
|
||||||
|
|
||||||
✅ **Proactive hunting** - Scours git servers every 2 min (backoff), finds data automatically
|
✅ **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
|
✅ **Smart timing** - 3min delay for user pushes, 500ms for synced events
|
||||||
✅ **30min expiry** - Auto-cleanup of events when data never arrives
|
✅ **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
|
✅ **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:
|
With throttling:
|
||||||
|
|
||||||
- ✅ Respect server capacity (5 concurrent max per domain)
|
- ✅ 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
|
- ✅ Fair access for all clients
|
||||||
|
|
||||||
### Two-Level Limits
|
### 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**:
|
**Example**:
|
||||||
|
|
||||||
@@ -292,23 +292,20 @@ fetch-3 completes → in-flight: 4
|
|||||||
Status: HAS CAPACITY (process next queued identifier)
|
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.
|
||||||
|
|
||||||
```
|
**Implementation**: [`git_budget.rs`](../../src/purgatory/sync/git_budget.rs).
|
||||||
t=0s: Request 1 → request_times: [0s]
|
The existing [pass scheduler](../../src/purgatory/sync/throttle.rs) also retains
|
||||||
t=1s: Request 2 → request_times: [0s, 1s]
|
its 60-pass/minute admission ceiling.
|
||||||
...
|
|
||||||
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)
|
|
||||||
|
|
||||||
### Round-Robin Fairness
|
### 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
|
3. **Residual OIDs one at a time** — needed OIDs that were neither
|
||||||
advertised nor discovered as ancestors of the fetched tips are
|
advertised nor discovered as ancestors of the fetched tips are
|
||||||
requested individually. Some servers refuse arbitrary-SHA1 wants, so a
|
requested individually. Some servers refuse arbitrary-SHA1 wants, so a
|
||||||
per-OID failure costs exactly one small round trip and never aborts the
|
missing-object rejection does not fail other OIDs. Speculative requests
|
||||||
rest of the pass.
|
are bounded as described below.
|
||||||
|
|
||||||
Missing residual OIDs are memoized per clone URL against a fingerprint of
|
Missing residual OIDs are memoized per clone URL against a fingerprint of
|
||||||
the sorted advertised OID set. Later passes skip those OIDs while the
|
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
|
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.
|
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
|
All three phases run beneath the pass's per-domain permit. That scheduler
|
||||||
comparison into its own helper in future must retain that permit; `ls-remote`
|
retains its five-concurrent-pass and 60-pass/minute limits. A separate shared
|
||||||
is outbound work against the same Git server, not an unaccounted preflight.
|
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
|
### 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://server1.com/repo.git", Ok(vec!["oid1"]))
|
||||||
.with_fetch_result("https://server2.com/repo.git", Ok(vec!["oid2"]));
|
.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;
|
let complete = sync_identifier(&mock, "repo", &throttle).await;
|
||||||
|
|
||||||
assert!(complete); // Both OIDs fetched
|
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 |
|
| Setting | CLI Flag | Environment Variable | Default | Description |
|
||||||
| ----------------------- | -------- | -------------------- | ------- | ---------------------------------------------------- |
|
| ----------------------- | -------- | -------------------- | ------- | ---------------------------------------------------- |
|
||||||
| Domain concurrent limit | (future) | (future) | `5` | Max concurrent requests per domain |
|
| 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) |
|
| 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) |
|
| 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) |
|
| 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:
|
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)
|
(one ls-remote comparison plus fetches per pass)
|
||||||
- `ngit_purgatory_git_fetch_oids_total{kind}` - Per-pass OID outcomes
|
- `ngit_purgatory_git_fetch_oids_total{kind}` - Per-pass OID outcomes
|
||||||
(`advertised_tip`, `residual_attempted`, `residual_missing`,
|
(`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
|
### 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`
|
**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.
|
**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
|
✅ **Batches intelligently** - Groups events by identifier for efficient fetching
|
||||||
✅ **Retries smartly** - Exponential backoff with fresh start on new events
|
✅ **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
|
✅ **Times strategically** - 3min for user events, 500ms for synced events
|
||||||
✅ **Expires responsibly** - 30min auto-cleanup prevents memory leaks
|
✅ **Expires responsibly** - 30min auto-cleanup prevents memory leaks
|
||||||
✅ **Soft-expires announcements** - Bare repo deleted at 30min, event retained 24h for revival
|
✅ **Soft-expires announcements** - Bare repo deleted at 30min, event retained 24h for revival
|
||||||
|
|||||||
+19
-8
@@ -520,14 +520,25 @@ impl FamilyRepairSource for crate::purgatory::sync::RealSyncContext {
|
|||||||
url: &str,
|
url: &str,
|
||||||
oids: &[String],
|
oids: &[String],
|
||||||
) -> Result<Vec<String>> {
|
) -> Result<Vec<String>> {
|
||||||
crate::purgatory::sync::SyncContext::fetch_oids_with_role(
|
loop {
|
||||||
self,
|
let result = crate::purgatory::sync::SyncContext::fetch_oids_with_role(
|
||||||
target_view,
|
self,
|
||||||
url,
|
target_view,
|
||||||
oids,
|
url,
|
||||||
crate::purgatory::sync::GitFetchRole::Integrity,
|
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;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -261,6 +261,15 @@ lazy_static! {
|
|||||||
REGISTRY.register(Box::new(metric.clone())).expect("register metric");
|
REGISTRY.register(Box::new(metric.clone())).expect("register metric");
|
||||||
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 = {
|
static ref PURGATORY_GIT_SUBPROCESS_TOTAL: CounterVec = {
|
||||||
let metric = CounterVec::new(
|
let metric = CounterVec::new(
|
||||||
Opts::new(
|
Opts::new(
|
||||||
@@ -358,6 +367,12 @@ pub struct PurgatoryGitSubprocessGuard {
|
|||||||
outcome: &'static str,
|
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 {
|
impl PurgatoryGitSubprocessGuard {
|
||||||
pub fn finish(&mut self, success: bool) {
|
pub fn finish(&mut self, success: bool) {
|
||||||
self.outcome = if success { "success" } else { "failure" };
|
self.outcome = if success { "success" } else { "failure" };
|
||||||
|
|||||||
+174
-19
@@ -240,6 +240,10 @@ use crate::purgatory::Purgatory;
|
|||||||
use crate::sync::naughty_list::NaughtyListTracker;
|
use crate::sync::naughty_list::NaughtyListTracker;
|
||||||
|
|
||||||
use super::functions::extract_domain;
|
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.
|
/// Real implementation of `SyncContext` that connects to actual systems.
|
||||||
///
|
///
|
||||||
@@ -274,6 +278,9 @@ pub struct RealSyncContext {
|
|||||||
/// advertised ref tips.
|
/// advertised ref tips.
|
||||||
miss_memo: Arc<Mutex<HashMap<String, RemoteMissMemo>>>,
|
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 target policy applied before every event-directed git fetch
|
||||||
outbound_policy: OutboundTargetPolicy,
|
outbound_policy: OutboundTargetPolicy,
|
||||||
|
|
||||||
@@ -323,6 +330,7 @@ impl RealSyncContext {
|
|||||||
write_policy,
|
write_policy,
|
||||||
git_naughty_list,
|
git_naughty_list,
|
||||||
miss_memo: Arc::new(Mutex::new(HashMap::new())),
|
miss_memo: Arc::new(Mutex::new(HashMap::new())),
|
||||||
|
git_budget: GitCommandBudget::default(),
|
||||||
outbound_policy,
|
outbound_policy,
|
||||||
grasp08_peers,
|
grasp08_peers,
|
||||||
credential_keys,
|
credential_keys,
|
||||||
@@ -676,6 +684,44 @@ async fn drain_git_stream<R: AsyncRead + Unpin>(
|
|||||||
Ok(captured)
|
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(
|
async fn run_observed_git_command(
|
||||||
command: Command,
|
command: Command,
|
||||||
domain: &str,
|
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);
|
crate::metrics::record_purgatory_git_subprocess_output(operation, role, "stderr", stderr_bytes);
|
||||||
if stalled {
|
if stalled {
|
||||||
metric.finish_with_outcome("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 {
|
} else {
|
||||||
metric.finish(status.success());
|
metric.finish(status.success());
|
||||||
}
|
}
|
||||||
@@ -1151,8 +1199,21 @@ impl SyncContext for RealSyncContext {
|
|||||||
// advertisement tells us up front which OIDs the remote can
|
// advertisement tells us up front which OIDs the remote can
|
||||||
// serve — instead of discovering each missing one through a
|
// serve — instead of discovering each missing one through a
|
||||||
// failed "not our ref" upload-pack round trip.
|
// 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 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(
|
hardened_git_command(
|
||||||
&repo_path,
|
&repo_path,
|
||||||
resolve_pin.as_deref(),
|
resolve_pin.as_deref(),
|
||||||
@@ -1189,13 +1250,7 @@ impl SyncContext for RealSyncContext {
|
|||||||
stderr
|
stderr
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => return Err(e.context(format!("git ls-remote command error for {url}"))),
|
||||||
return Err(anyhow::anyhow!(
|
|
||||||
"git ls-remote command error for {}: {}",
|
|
||||||
url,
|
|
||||||
e
|
|
||||||
))
|
|
||||||
}
|
|
||||||
};
|
};
|
||||||
let advert_fingerprint = advertised_oids_fingerprint(&advertised);
|
let advert_fingerprint = advertised_oids_fingerprint(&advertised);
|
||||||
let memoized_missing = {
|
let memoized_missing = {
|
||||||
@@ -1224,7 +1279,9 @@ impl SyncContext for RealSyncContext {
|
|||||||
];
|
];
|
||||||
args.extend(advertised_tips.iter().cloned());
|
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(
|
hardened_git_command(
|
||||||
&repo_path,
|
&repo_path,
|
||||||
resolve_pin.as_deref(),
|
resolve_pin.as_deref(),
|
||||||
@@ -1260,13 +1317,7 @@ impl SyncContext for RealSyncContext {
|
|||||||
return Err(anyhow::anyhow!("git fetch failed for {}: {}", url, stderr));
|
return Err(anyhow::anyhow!("git fetch failed for {}: {}", url, stderr));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => return Err(e.context(format!("git fetch command error for {url}"))),
|
||||||
return Err(anyhow::anyhow!(
|
|
||||||
"git fetch command error for {}: {}",
|
|
||||||
url,
|
|
||||||
e
|
|
||||||
))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1275,6 +1326,7 @@ impl SyncContext for RealSyncContext {
|
|||||||
// time. Arbitrary-SHA1 wants may still be refused by some
|
// time. Arbitrary-SHA1 wants may still be refused by some
|
||||||
// servers, so a per-OID failure must cost exactly one small
|
// servers, so a per-OID failure must cost exactly one small
|
||||||
// round trip and must not abort the rest.
|
// round trip and must not abort the rest.
|
||||||
|
let mut deferred = None;
|
||||||
let mut residual_attempted = 0usize;
|
let mut residual_attempted = 0usize;
|
||||||
let mut residual_missing = 0usize;
|
let mut residual_missing = 0usize;
|
||||||
let mut residual_skipped = 0usize;
|
let mut residual_skipped = 0usize;
|
||||||
@@ -1286,7 +1338,10 @@ impl SyncContext for RealSyncContext {
|
|||||||
residual_skipped += 1;
|
residual_skipped += 1;
|
||||||
continue;
|
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![
|
let args = vec![
|
||||||
"fetch".to_string(),
|
"fetch".to_string(),
|
||||||
@@ -1296,7 +1351,9 @@ impl SyncContext for RealSyncContext {
|
|||||||
url.clone(),
|
url.clone(),
|
||||||
oid.clone(),
|
oid.clone(),
|
||||||
];
|
];
|
||||||
match run_observed_git_command(
|
let result = run_budgeted_git_command(
|
||||||
|
&self.git_budget,
|
||||||
|
reservation.as_mut(),
|
||||||
hardened_git_command(
|
hardened_git_command(
|
||||||
&repo_path,
|
&repo_path,
|
||||||
resolve_pin.as_deref(),
|
resolve_pin.as_deref(),
|
||||||
@@ -1309,8 +1366,16 @@ impl SyncContext for RealSyncContext {
|
|||||||
"fetch_residual",
|
"fetch_residual",
|
||||||
role,
|
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.status.success() => {}
|
||||||
Ok(result) => {
|
Ok(result) => {
|
||||||
if result.stderr_truncated {
|
if result.stderr_truncated {
|
||||||
@@ -1397,6 +1462,11 @@ impl SyncContext for RealSyncContext {
|
|||||||
fetched.len(),
|
fetched.len(),
|
||||||
);
|
);
|
||||||
|
|
||||||
|
if role == GitFetchRole::Integrity {
|
||||||
|
if let Some(error) = deferred {
|
||||||
|
return Err(error);
|
||||||
|
}
|
||||||
|
}
|
||||||
Ok(fetched)
|
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]
|
#[tokio::test]
|
||||||
async fn stream_drain_bounds_capture_but_counts_all_activity() {
|
async fn stream_drain_bounds_capture_but_counts_all_activity() {
|
||||||
let (mut writer, reader) = tokio::io::duplex(256);
|
let (mut writer, reader) = tokio::io::duplex(256);
|
||||||
|
|||||||
@@ -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"
|
||||||
|
));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -9,6 +9,7 @@
|
|||||||
|
|
||||||
mod context;
|
mod context;
|
||||||
mod functions;
|
mod functions;
|
||||||
|
mod git_budget;
|
||||||
mod r#loop;
|
mod r#loop;
|
||||||
mod queue;
|
mod queue;
|
||||||
mod throttle;
|
mod throttle;
|
||||||
@@ -18,6 +19,7 @@ pub use functions::{
|
|||||||
get_throttled_domains_with_untried_urls, sync_identifier, sync_identifier_from_url,
|
get_throttled_domains_with_untried_urls, sync_identifier, sync_identifier_from_url,
|
||||||
sync_identifier_next_url, ThrottledDomainInfo,
|
sync_identifier_next_url, ThrottledDomainInfo,
|
||||||
};
|
};
|
||||||
|
pub(crate) use git_budget::GitBudgetDeferred;
|
||||||
pub use queue::SyncQueueEntry;
|
pub use queue::SyncQueueEntry;
|
||||||
pub use throttle::{DomainThrottle, ThrottleManager};
|
pub use throttle::{DomainThrottle, ThrottleManager};
|
||||||
|
|
||||||
|
|||||||
@@ -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
|
//! This module provides per-domain throttling to prevent overwhelming remote
|
||||||
//! git servers during purgatory sync operations. Each domain has:
|
//! git servers during purgatory sync operations. Each domain has:
|
||||||
|
|||||||
+3
-2
@@ -480,8 +480,9 @@ impl RelayServer {
|
|||||||
));
|
));
|
||||||
info!("Git storage and authorization-integrity worker started");
|
info!("Git storage and authorization-integrity worker started");
|
||||||
|
|
||||||
// Create throttle manager for rate limiting remote git servers
|
// Bound scheduled fetch passes (5 concurrent, 60 starts/minute/domain).
|
||||||
// Default: 5 concurrent requests per domain, 60 requests per minute per domain
|
// RealSyncContext separately budgets every outbound Git command,
|
||||||
|
// including commands issued by integrity repair.
|
||||||
let throttle_manager = Arc::new(ThrottleManager::new(5, 60));
|
let throttle_manager = Arc::new(ThrottleManager::new(5, 60));
|
||||||
throttle_manager.set_context(sync_ctx.clone());
|
throttle_manager.set_context(sync_ctx.clone());
|
||||||
throttle_manager.set_git_naughty_list(git_naughty_list.clone());
|
throttle_manager.set_git_naughty_list(git_naughty_list.clone());
|
||||||
|
|||||||
@@ -22,7 +22,7 @@
|
|||||||
//! The scenario drives the real purgatory sync path end to end: a genuine
|
//! 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
|
//! 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
|
//! 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
|
//! 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.
|
//! 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
|
/// Declared ref tips that exist on no reachable server, mirroring the
|
||||||
/// `market` repository's unfetchable state event.
|
/// `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`.
|
/// 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 {
|
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
|
/// Wait for the fetch summary, which is emitted after residual admission and
|
||||||
/// exchanges and the count has been stable for `stable_for`.
|
/// object retention finish. Object arrival alone precedes those operations.
|
||||||
async fn wait_for_upload_pack_quiescence(
|
async fn wait_for_completed_passes(relay: &TestRelay, minimum: usize) -> Vec<String> {
|
||||||
proxy: &UploadPackCountingProxy,
|
let deadline = tokio::time::Instant::now() + Duration::from_secs(60);
|
||||||
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();
|
|
||||||
loop {
|
loop {
|
||||||
let total = proxy.exchanges().len();
|
let log = std::fs::read_to_string(relay.log_path()).unwrap_or_default();
|
||||||
if total != last_total {
|
let passes: Vec<String> = log
|
||||||
last_total = total;
|
.lines()
|
||||||
stable_since = tokio::time::Instant::now();
|
.filter(|line| line.contains("Purgatory git fetch pass complete"))
|
||||||
} else if total >= min_total && stable_since.elapsed() >= stable_for {
|
.map(str::to_owned)
|
||||||
return true;
|
.collect();
|
||||||
|
if passes.len() >= minimum {
|
||||||
|
return passes;
|
||||||
}
|
}
|
||||||
if tokio::time::Instant::now() >= end {
|
assert!(
|
||||||
return false;
|
tokio::time::Instant::now() < deadline,
|
||||||
}
|
"fetch pass should complete"
|
||||||
tokio::time::sleep(Duration::from_millis(200)).await;
|
);
|
||||||
|
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.
|
/// 2. A counting proxy fronts the source's git smart-HTTP endpoint.
|
||||||
/// 3. The relay under test bootstrap-syncs from a MockRelay serving a
|
/// 3. The relay under test bootstrap-syncs from a MockRelay serving a
|
||||||
/// later announcement whose clone tag points at the proxy, plus a state
|
/// 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
|
/// Both events sit in purgatory; the purgatory sync loop fetches
|
||||||
/// through the proxy.
|
/// through the proxy.
|
||||||
/// 4. The two real tips must arrive, and the outbound request shape must
|
/// 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).
|
/// single missing object must never fail a batch).
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
|
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
|
// 1. Source relay + counting proxy in front of its git endpoint, and
|
||||||
// the MockRelay that will carry the events for the relay under test.
|
// the MockRelay that will carry the events for the relay under test.
|
||||||
let source = TestRelay::start().await;
|
let source = TestRelay::start().await;
|
||||||
@@ -208,10 +213,10 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
|
|||||||
.finalize(&keys)
|
.finalize(&keys)
|
||||||
.expect("sign syncing announcement");
|
.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}"))
|
.map(|index| format!("beef{index:036x}"))
|
||||||
.collect();
|
.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}"))
|
.map(|index| format!("missing-{index}"))
|
||||||
.collect();
|
.collect();
|
||||||
let mut branches: Vec<(&str, &str)> =
|
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");
|
.expect("send state event to mock relay");
|
||||||
|
|
||||||
// Negentropy is disabled because MockRelay does not support NIP-77.
|
// 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,
|
syncing_reservation,
|
||||||
Some(mock.url().to_string()),
|
Some(mock.url().to_string()),
|
||||||
true,
|
true,
|
||||||
|
archive,
|
||||||
|
archive,
|
||||||
)
|
)
|
||||||
.await;
|
.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
|
// 6. Let the sync pass finish so requests aimed at the missing tips
|
||||||
// (however the client shapes them) are all recorded.
|
// (however the client shapes them) are all recorded.
|
||||||
wait_for_upload_pack_quiescence(&proxy, 1, Duration::from_secs(3), Duration::from_secs(60))
|
let first_passes = wait_for_completed_passes(&syncing, 1).await;
|
||||||
.await;
|
|
||||||
|
|
||||||
// 7. Regression assertions on the outbound request shape.
|
// 7. Regression assertions on the outbound request shape.
|
||||||
let exchanges = proxy.exchanges();
|
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
|
// 9. A fresh replaceable state event resets the queue backoff and
|
||||||
// deterministically triggers another pass. With an unchanged remote
|
// 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
|
let failed_after_first_pass = proxy
|
||||||
.exchanges()
|
.exchanges()
|
||||||
.iter()
|
.iter()
|
||||||
@@ -343,32 +350,54 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() {
|
|||||||
.await
|
.await
|
||||||
.expect("send fresh state event to mock relay");
|
.expect("send fresh state event to mock relay");
|
||||||
|
|
||||||
let second_pass_deadline = tokio::time::Instant::now() + Duration::from_secs(60);
|
let passes = wait_for_completed_passes(&syncing, first_passes.len() + 1).await;
|
||||||
while proxy.info_refs_count() <= info_refs_before_second_pass {
|
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!(
|
assert!(
|
||||||
tokio::time::Instant::now() < second_pass_deadline,
|
attempted <= 2,
|
||||||
"fresh state event should trigger a second ls-remote"
|
"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
|
let failed_after_second_pass = proxy
|
||||||
.exchanges()
|
.exchanges()
|
||||||
.iter()
|
.iter()
|
||||||
.filter(|exchange| !exchange.ok)
|
.filter(|exchange| !exchange.ok)
|
||||||
.count();
|
.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 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!(
|
assert_eq!(
|
||||||
failed_after_second_pass, failed_after_first_pass,
|
unique.len(),
|
||||||
"unchanged advertisement must produce no new failed upload-pack exchanges"
|
failed_oids.len(),
|
||||||
|
"memoized misses must not be requested again"
|
||||||
);
|
);
|
||||||
|
|
||||||
source_client.disconnect().await;
|
source_client.disconnect().await;
|
||||||
|
|||||||
Reference in New Issue
Block a user