From cd0f0d94b7bb5b55eab1464e711c2f4786ab9937 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 15:44:35 -0700 Subject: [PATCH 01/12] Reduce Git fetch contention and bound per-account waiting --- docs/performance-plan.md | 46 +++++++++++++++-- scripts/benchmark_repositories.py | 31 ++++++++++-- scripts/test_benchmark_repositories.py | 44 ++++++++++++++++ src/admission.rs | 69 +++++++++++++++++++++++--- src/git_cache.rs | 52 ++++++++++++++----- src/git_cache/tests.rs | 36 ++++++++++++++ src/git_gateway/fetch.rs | 16 +++--- tests/multi_server/compatibility.rs | 39 +++++++++++++++ 8 files changed, 298 insertions(+), 35 deletions(-) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index de890c6..28a2e6b 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -71,6 +71,12 @@ branch adds its blob and leaves the deleted branch absent. This runs for v0/v2 over both transports, and establishes reduced preparation, not a latency target. Preparation logs report newly hydrated blob counts/bytes and elapsed time. +Concurrent fetches of one repository now share the immutable object cache +without holding its gateway mutex for the entire preparation. Verified loose +objects publish atomically; writers for the same object ID converge while +different IDs may hydrate in parallel. A cold concurrent stock-Git clone test +checks complete packs and `git fsck`. This removes one serialization point; +it does not establish hot-repository throughput or a 10,000-repository result. The native walk still traverses requested history on warm requests; native Git's own traversal memory is not bounded by the Rust batch size. Structural hydration pages certified edges and retains a visited set proportional to the requested @@ -175,10 +181,12 @@ permits it. object metadata using an insertion cursor and reuses verified immutable bodies from a repository-scoped cache. Each native push/merge has a private writable generation; successful publication precedes hydration of its new objects into the shared cache. -- Eight shared heavy-request slots, a four-slot ceiling per account, and the - bounded Linux profile establish containment. They do not establish latency, throughput or thousands of active - repositories. The profile's tmpfs charges cache bytes to memory; density - qualification needs a separate bounded NVMe-backed profile. +- Eight shared heavy-request slots and four active and four pending slots per + account prevent one account from filling every transfer wait position. The + bounded Linux profile establishes containment, not latency, throughput or + thousands of active repositories. Its tmpfs charges cache bytes to memory; + density qualification needs a separate bounded NVMe-backed profile and + mixed-account latency measurements. - The Linux profile sets both process descriptor limits to 16,384. Its checker requires eight per admitted Repository/Directory Cell plus 1,024 headroom; this is a reservation check, not measured aggregate process capacity. @@ -397,6 +405,30 @@ empty. This is a density smoke corpus, not realistic large-history qualification An interrupted seed leaves an explicitly incomplete manifest and cannot be used as a complete corpus. Names have a unique prefix; the manifest records them. +For a multi-gateway read run, seed the documented 10,000-repository corpus, +start two nodes against the same deployment, and pass each ingress to `run`: + +```bash +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --manifest /path/to/canopy-corpus.json \ + seed --repositories 10000 --populated 100 \ + --work-dir /path/to/canopy-seed +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-corpus.json \ + run --active-repositories 1000 --distribution uniform \ + --operation refs --rate 20 --duration 120 --concurrency 32 \ + --output /path/to/canopy-two-gateway-refs.json +``` + +The run assigns scheduled requests to ingresses by sequence number and records +each ingress's outcomes and latency percentiles. Compare the same workload with one and two ingresses +and collect node, object-store, and client resource metrics separately. This +driver measures metadata or Git v2 discovery; it does not measure push, pack +generation, full fetch, or LFS transfer throughput. + Run `--operation refs` for Git v2 discovery or `--distribution skewed` for 90% of requests to the selected working set's first tenth. A fixed seed determines working-set selection and offered arrivals. The driver never calls a working set @@ -1391,6 +1423,12 @@ Retry-After response, other-account progress and disconnect recovery; a unit test saturates pending admission, exercises another account and cancels a waiter. State is bounded by active/pending work, not repository or account history. +The current admission implementation also caps pending positions at four per +account, leaving shared wait positions for other accounts. This is a scheduling +bound, not a throughput result. Measure it under mixed-account Git and LFS +traffic before treating it as a fairness guarantee at the 10,000-repository +reference target. + Real-provider run `canopy-transfer-fairness-013e78122a08` used release binary `7e5f0a173cebd77aa65cbbe646f2566b75f462875653381f72a3795ddab48b57` (base `549f775` plus this patch), macOS arm64 and RustFS diff --git a/scripts/benchmark_repositories.py b/scripts/benchmark_repositories.py index 5c0ea1b..7e9d4a4 100644 --- a/scripts/benchmark_repositories.py +++ b/scripts/benchmark_repositories.py @@ -184,6 +184,7 @@ def percentiles(values): def measure(args, client, token): del token + clients = client if isinstance(client, list) else [client] manifest = corpus(args.manifest) if args.active_repositories > len(manifest["repositories"]): raise ValueError("active repository count exceeds corpus") @@ -195,6 +196,9 @@ def measure(args, client, token): if total > 1_000_000: raise ValueError("one run is limited to one million scheduled arrivals") counts, latencies, service_times, dispatch_times = Counter(), [], [], [] + ingress_counts = [Counter() for _ in clients] + ingress_latencies = [[] for _ in clients] + ingress_service_times = [[] for _ in clients] slots = threading.BoundedSemaphore(args.concurrency) lock = threading.Lock() started = time.monotonic() @@ -206,30 +210,35 @@ def measure(args, client, token): def record(sample): with lock: counts[sample["result"]] += 1 + ingress_counts[sample["ingress_index"]][sample["result"]] += 1 if sample["elapsed_ms"] is not None: latencies.append(sample["elapsed_ms"]) service_times.append(sample["service_ms"]) dispatch_times.append(sample["dispatch_delay_ms"]) + ingress_latencies[sample["ingress_index"]].append(sample["elapsed_ms"]) + ingress_service_times[sample["ingress_index"]].append(sample["service_ms"]) samples.write(json.dumps(sample) + "\n") def execute(sequence, entry, scheduled): request_id = str(uuid.uuid4()) + ingress_index = sequence % len(clients) + ingress = clients[ingress_index] dispatched = time.monotonic() result = "transport_error" try: if args.operation == "metadata": - status, body = client.request(f"/api/repositories/{entry['name']}", request_id=request_id) + status, body = ingress.request(f"/api/repositories/{entry['name']}", request_id=request_id) decoded = json.loads(body) if status == 200 else None valid = isinstance(decoded, dict) and decoded.get("repository_id") == entry["repository_id"] else: - status, body = client.request(f"/{entry['owner']}/{entry['name']}.git/info/refs?service=git-upload-pack", git=True, request_id=request_id) + status, body = ingress.request(f"/{entry['owner']}/{entry['name']}.git/info/refs?service=git-upload-pack", git=True, request_id=request_id) valid = status == 200 and body.startswith(b"000eversion 2\n") result = "ok" if valid else (f"http_{status}" if status != 200 else "invalid_response") except (OSError, ValueError, KeyError, http.client.HTTPException): pass finally: finished = time.monotonic() - record({"sequence": sequence, "request_id": request_id, "repository_id": entry["repository_id"], "result": result, + record({"sequence": sequence, "request_id": request_id, "repository_id": entry["repository_id"], "ingress_index": ingress_index, "result": result, "elapsed_ms": (finished - scheduled) * 1000, "service_ms": (finished - dispatched) * 1000, "dispatch_delay_ms": (dispatched - scheduled) * 1000}) @@ -247,6 +256,7 @@ def execute(sequence, entry, scheduled): executor.submit(execute, sequence, entry, scheduled) else: record({"sequence": sequence, "repository_id": entry["repository_id"], + "ingress_index": sequence % len(clients), "result": "driver_busy", "elapsed_ms": None}) elapsed = time.monotonic() - started result = {"version": 1, "started_at_utc": started_at, @@ -257,6 +267,10 @@ def execute(sequence, entry, scheduled): "elapsed_including_drain_seconds": round(elapsed, 3), "concurrency": args.concurrency, "request_timeout_seconds": args.timeout, "scheduled": total, "outcomes": dict(counts), + "ingresses": [{"index": index, "outcomes": dict(outcomes), + "scheduled_latency_ms": percentiles(ingress_latencies[index]), + "service_ms": percentiles(ingress_service_times[index])} + for index, outcomes in enumerate(ingress_counts)], "failed_arrivals": total - counts["ok"], "scheduled_latency_ms": percentiles(latencies), "service_ms": percentiles(service_times), "dispatch_delay_ms": percentiles(dispatch_times), @@ -278,6 +292,8 @@ def positive(value): def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--base-url", required=True) + parser.add_argument("--additional-base-url", action="append", default=[], + help="additional gateway ingress for round-robin run requests") parser.add_argument("--manifest", type=Path, required=True) parser.add_argument("--timeout", type=positive, default=10) parser.add_argument("--seed", type=int, default=20260926) @@ -302,12 +318,17 @@ def main(): parser.error("CANOPY_GIT_TOKEN is required") if args.command == "run" and args.concurrency > 256: parser.error("driver concurrency is limited to 256") - client = Client(args.base_url, token, args.timeout) + if args.command != "run" and args.additional_base_url: + parser.error("additional gateways are supported only for run") + clients = [Client(base, token, args.timeout) + for base in [args.base_url, *args.additional_base_url]] try: + client = clients if args.command == "run" else clients[0] result = {"seed": seed, "verify": verify, "run": measure}[args.command](args, client, token) print(json.dumps(result, indent=2)) finally: - client.close() + for client in clients: + client.close() if args.command == "run" and result["failed_arrivals"]: raise SystemExit(1) diff --git a/scripts/test_benchmark_repositories.py b/scripts/test_benchmark_repositories.py index e6f2f1b..9c3a0f0 100644 --- a/scripts/test_benchmark_repositories.py +++ b/scripts/test_benchmark_repositories.py @@ -78,6 +78,50 @@ def test_overload_has_no_hidden_retries_or_missing_arrivals(self): server.server_close() thread.join() + def test_multiple_ingresses_report_each_gateway_without_losing_arrivals(self): + servers = [ThreadingHTTPServer(("127.0.0.1", 0), Handler) for _ in range(2)] + threads = [] + clients = [] + for server in servers: + server.guard = threading.Lock() + server.requests = server.active = server.peak = 0 + server.request_ids = [] + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + threads.append(thread) + clients.append(benchmark.Client(f"http://127.0.0.1:{server.server_port}", + "fixture-token", 2)) + try: + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + manifest = root / "manifest.json" + manifest.write_text(json.dumps({"version": 1, "complete": True, + "requested_repositories": 1, "repositories": [{"name": "fixture", "owner": "canopy", + "repository_id": str(uuid.uuid4()), "commit": None}]})) + args = SimpleNamespace(manifest=manifest, active_repositories=1, seed=42, + duration=1, rate=20, concurrency=16, distribution="uniform", + operation="metadata", timeout=2, output=root / "report.json") + report = benchmark.measure(args, clients, None) + samples = [json.loads(line) for line in + (root / "report.samples.jsonl").read_text().splitlines()] + self.assertEqual(sorted(sample["sequence"] for sample in samples), list(range(20))) + self.assertEqual({sample["ingress_index"] for sample in samples}, {0, 1}) + self.assertEqual(sum(server.requests for server in servers), + report["outcomes"].get("http_503", 0)) + for index, server in enumerate(servers): + self.assertEqual(server.requests, + report["ingresses"][index]["outcomes"].get("http_503", 0)) + self.assertGreater(server.requests, 0) + self.assertIsNotNone(report["ingresses"][index]["service_ms"]["p95"]) + finally: + for client in clients: + client.close() + for server in servers: + server.shutdown() + server.server_close() + for thread in threads: + thread.join() + if __name__ == "__main__": unittest.main() diff --git a/src/admission.rs b/src/admission.rs index 1bedff7..9809940 100644 --- a/src/admission.rs +++ b/src/admission.rs @@ -22,6 +22,7 @@ pub(crate) struct AccountAdmission { total_capacity: &'static str, account_capacity: &'static str, accounts: Mutex, Weak>>, + waiting_accounts: Mutex, Weak>>, } impl AccountAdmission { @@ -37,6 +38,7 @@ impl AccountAdmission { total_capacity, account_capacity, accounts: Mutex::new(HashMap::new()), + waiting_accounts: Mutex::new(HashMap::new()), } } @@ -57,8 +59,13 @@ impl AccountAdmission { } pub(crate) async fn wait(&self, actor: ReadIdentity<'_>) -> Result { - // Bound retained HTTP requests before waiting. Acquire the account first: - // a busy account must not reserve node slots needed by another account. + // Bound each account's waiters before the shared queue. A busy account + // cannot fill every pending position while another account has work. + let _account_waiting = self + .waiting_account(actor) + .await + .try_acquire_owned() + .map_err(|_| Error::Capacity(self.account_capacity))?; let _waiting = self .waiting .try_acquire() @@ -80,19 +87,31 @@ impl AccountAdmission { } async fn account(&self, actor: ReadIdentity<'_>) -> Arc { + Self::semaphore(&self.accounts, actor, self.account_limit).await + } + + async fn waiting_account(&self, actor: ReadIdentity<'_>) -> Arc { + Self::semaphore(&self.waiting_accounts, actor, self.account_limit).await + } + + async fn semaphore( + entries: &Mutex, Weak>>, + actor: ReadIdentity<'_>, + limit: usize, + ) -> Arc { let account = match actor { ReadIdentity::Account(account) => Some(account.to_owned()), ReadIdentity::Anonymous => None, }; { - let mut accounts = self.accounts.lock().await; + let mut accounts = entries.lock().await; // Permits retain their semaphore through detached ownership work. // Active or waiting admission bounds this map; expired accounts need no state. accounts.retain(|_, semaphore| semaphore.strong_count() != 0); if let Some(semaphore) = accounts.get(&account).and_then(Weak::upgrade) { semaphore } else { - let semaphore = Arc::new(Semaphore::new(self.account_limit)); + let semaphore = Arc::new(Semaphore::new(limit)); accounts.insert(account, Arc::downgrade(&semaphore)); semaphore } @@ -140,10 +159,8 @@ mod tests { let actor = ReadIdentity::Account("busy"); let held = admission.acquire(actor).await.unwrap(); let mut cancelled = Box::pin(admission.wait(actor)); - let mut pending = Box::pin(admission.wait(actor)); poll_fn(|cx| { assert!(cancelled.as_mut().poll(cx).is_pending()); - assert!(pending.as_mut().poll(cx).is_pending()); Poll::Ready(()) }) .await; @@ -153,7 +170,13 @@ mod tests { .await .unwrap(); drop(cancelled); - assert_eq!(admission.waiting.available_permits(), 1); + assert_eq!(admission.waiting.available_permits(), 2); + let mut pending = Box::pin(admission.wait(actor)); + poll_fn(|cx| { + assert!(pending.as_mut().poll(cx).is_pending()); + Poll::Ready(()) + }) + .await; drop(held); let admitted = pending.await.unwrap(); assert_eq!(admission.total.available_permits(), 0); @@ -163,6 +186,38 @@ mod tests { assert_eq!(admission.total.available_permits(), 2); } + #[tokio::test] + async fn one_account_cannot_fill_the_shared_waiting_queue() { + use std::{ + future::{Future, poll_fn}, + task::Poll, + }; + let admission = AccountAdmission::new(8, "total", "account"); + let busy = ReadIdentity::Account("busy"); + let other = ReadIdentity::Account("other"); + let mut busy_active = Vec::new(); + for _ in 0..4 { + busy_active.push(admission.acquire(busy).await.unwrap()); + } + let mut busy_waiters: Vec<_> = (0..4).map(|_| Box::pin(admission.wait(busy))).collect(); + for waiter in &mut busy_waiters { + poll_fn(|cx| { + assert!(waiter.as_mut().poll(cx).is_pending()); + Poll::Ready(()) + }) + .await; + } + assert_eq!(admission.waiting.available_permits(), 4); + assert!(admission.wait(busy).await.is_err()); + let other_permit = admission.wait(other).await.unwrap(); + assert_eq!(admission.waiting.available_permits(), 4); + drop(other_permit); + drop(busy_waiters); + drop(busy_active); + assert_eq!(admission.waiting.available_permits(), 8); + assert_eq!(admission.total.available_permits(), 8); + } + #[tokio::test] async fn finished_accounts_do_not_accumulate_state() { let admission = AccountAdmission::new(32, "total", "account"); diff --git a/src/git_cache.rs b/src/git_cache.rs index 75b4660..997485d 100644 --- a/src/git_cache.rs +++ b/src/git_cache.rs @@ -10,6 +10,7 @@ use std::{ use cellule_ltx::{DiskBudget, DiskReservation, LtxError}; use flate2::{Compression, write::ZlibEncoder}; +use tokio::sync::Mutex; use crate::{ ObjectKind, RefExpectation, @@ -51,6 +52,9 @@ pub(crate) struct GitCache { directory: tempfile::TempDir, reservation: Option, objects: Option>, + // Only durable hydration writes this cache. Stripe by OID so concurrent + // fetches share a completed loose object without serializing all objects. + object_writes: [Arc>; 64], } impl GitCache { @@ -82,6 +86,7 @@ impl GitCache { directory: tempfile::Builder::new().prefix(CACHE_PREFIX).tempdir_in(fs::canonicalize(root)?)?, reservation: Some(budget.try_reserve(0)?), objects, + object_writes: std::array::from_fn(|_| Arc::new(Mutex::new(()))), }); for directory in ["objects/info", "objects/pack", "refs/heads", "refs/tags", "hooks"] { fs::create_dir_all(cache.git_dir().join(directory))?; @@ -147,17 +152,8 @@ impl GitCache { tokio::task::spawn_blocking(move || { let mut missing = Vec::new(); for oid in ids { - match fs::symlink_metadata(cache.object_path(oid)) { - Ok(metadata) if metadata.is_file() => {} - Err(error) if error.kind() == io::ErrorKind::NotFound => missing.push(oid), - Err(error) => return Err(CacheError::Io(error)), - _ => { - return Err(io::Error::new( - io::ErrorKind::InvalidData, - "invalid cached object", - ) - .into()); - } + if !cache.object_present(oid)? { + missing.push(oid); } } Ok(missing) @@ -173,6 +169,18 @@ impl GitCache { .join(&hex[2..]) } + fn object_present(&self, oid: crate::ObjectId) -> io::Result { + match fs::symlink_metadata(self.object_path(oid)) { + Ok(metadata) if metadata.is_file() => Ok(true), + Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(false), + Err(error) => Err(error), + _ => Err(io::Error::new( + io::ErrorKind::InvalidData, + "invalid cached object", + )), + } + } + fn object_writer( self: &Arc, oid: crate::ObjectId, @@ -237,11 +245,18 @@ impl GitCache { kind: ObjectKind, body: Vec, ) -> Result<(), CacheError> { + let write = Arc::clone(&self.object_writes[oid[0] as usize % self.object_writes.len()]) + .lock_owned() + .await; let cache = Arc::clone(self); tokio::task::spawn_blocking(move || { + let _write = write; if oid.format() != cache.object_format || object_id(oid.format(), kind, &body) != oid { return Err(io::Error::new(io::ErrorKind::InvalidData, "Git OID mismatch").into()); } + if cache.object_present(oid)? { + return Ok(()); + } let (mut encoder, temporary, destination) = cache.object_writer(oid)?; encoder.write_all(format!("{} {}\0", kind.git_name(), body.len()).as_bytes())?; encoder.write_all(&body)?; @@ -259,8 +274,12 @@ impl GitCache { mut reader: LargeBlobRead, ) -> Result<(), CacheError> { let reference = reader.reference(); + let write = + Arc::clone(&self.object_writes[reference.oid[0] as usize % self.object_writes.len()]) + .lock_owned() + .await; let cache = Arc::clone(self); - let (mut encoder, temporary, destination) = tokio::task::spawn_blocking(move || { + let pending = tokio::task::spawn_blocking(move || { if reference.oid.format() != cache.object_format { return Err(io::Error::new( io::ErrorKind::InvalidData, @@ -268,11 +287,17 @@ impl GitCache { ) .into()); } + if cache.object_present(reference.oid)? { + return Ok::<_, CacheError>(None); + } let (mut encoder, temporary, destination) = cache.object_writer(reference.oid)?; encoder.write_all(format!("blob {}\0", reference.size).as_bytes())?; - Ok::<_, CacheError>((encoder, temporary, destination)) + Ok::<_, CacheError>(Some((encoder, temporary, destination))) }) .await??; + let Some((mut encoder, temporary, destination)) = pending else { + return Ok(()); + }; while let Some(bytes) = reader.next().await? { // The writer owns the cache reservation until each compression worker // exits, including when its async waiter is canceled. @@ -283,6 +308,7 @@ impl GitCache { .await??; } tokio::task::spawn_blocking(move || { + let _write = write; drop(encoder.finish()?); temporary .persist_noclobber(destination) diff --git a/src/git_cache/tests.rs b/src/git_cache/tests.rs index deb1aab..cc77051 100644 --- a/src/git_cache/tests.rs +++ b/src/git_cache/tests.rs @@ -1,5 +1,41 @@ use super::*; +#[tokio::test] +async fn concurrent_hydration_publishes_each_object_once() -> Result<(), Box> +{ + let root = tempfile::TempDir::new()?; + let budget = DiskBudget::new(1 << 20); + let cache = GitCache::create( + root.path().into(), + budget.clone(), + "refs/heads/main", + crate::ObjectFormat::Sha1, + ) + .await?; + let body = b"shared by concurrent fetches".to_vec(); + let oid = object_id(crate::ObjectFormat::Sha1, ObjectKind::Blob, &body); + let mut workers = tokio::task::JoinSet::new(); + for _ in 0..32 { + let cache = Arc::clone(&cache); + let body = body.clone(); + workers.spawn(async move { cache.store_object(oid, ObjectKind::Blob, body).await }); + } + while let Some(result) = workers.join_next().await { + result??; + } + assert!(cache.missing_objects(vec![oid]).await?.is_empty()); + assert_eq!(budget.used(), tree_bytes(cache.root())?); + let output = tokio::process::Command::new("git") + .arg("--git-dir") + .arg(cache.git_dir()) + .args(["cat-file", "blob", &hex::encode(oid)]) + .output() + .await?; + assert!(output.status.success()); + assert_eq!(output.stdout, body); + Ok(()) +} + #[tokio::test] async fn hydrated_files_remain_charged_until_the_last_reader_releases_them() -> Result<(), Box> { diff --git a/src/git_gateway/fetch.rs b/src/git_gateway/fetch.rs index e6a13d2..13910f7 100644 --- a/src/git_gateway/fetch.rs +++ b/src/git_gateway/fetch.rs @@ -113,16 +113,20 @@ impl GitGateway { return Ok(()); } let started = std::time::Instant::now(); - let objects = self.objects.lock().await; - let shared = objects.as_ref().ok_or(GatewayError::MalformedCache)?; + // The shared object cache publishes loose objects atomically and + // coordinates duplicate OID writes. Hold the gateway lock only long + // enough to borrow it; slow fetches must not queue behind each other. + let shared = { + let objects = self.objects.lock().await; + Arc::clone(&objects.as_ref().ok_or(GatewayError::MalformedCache)?.cache) + }; let roots: Vec<_> = request.wants.iter().copied().collect(); - self.hydrate_selected(&shared.cache, request.wants).await?; + self.hydrate_selected(&shared, request.wants).await?; let unfiltered = request.filter.is_none(); // The certified Cell graph already names every reachable blob. A full // fetch can hydrate those bodies during the structural walk and avoid // a second native traversal over the same cold history. - self.hydrate_structure(&shared.cache, &roots, unfiltered) - .await?; + self.hydrate_structure(&shared, &roots, unfiltered).await?; if unfiltered || request.filter.as_deref() == Some("blob:none") { return Ok(()); } @@ -146,7 +150,7 @@ impl GitGateway { if ids.is_empty() { break; } - self.hydrate_objects(&shared.cache, ids, &mut stats).await?; + self.hydrate_objects(&shared, ids, &mut stats).await?; } walk.finish().await?; tracing::debug!( diff --git a/tests/multi_server/compatibility.rs b/tests/multi_server/compatibility.rs index f7b1875..b2fa926 100644 --- a/tests/multi_server/compatibility.rs +++ b/tests/multi_server/compatibility.rs @@ -77,6 +77,45 @@ async fn stock_git_history_refs_and_shallow_fetch_survive_fresh_disk_restore() - run_git(Some(&source), &["branch", branch]).await?; } run_git(Some(&source), &["-c", AUTH, "push", "--mirror", &url]).await?; + // All clones target one cold object cache at once. They must not race to + // publish duplicate loose objects or return incomplete packs. + let concurrent: Vec<_> = (0..4) + .map(|index| workspace.path().join(format!("concurrent-{index}"))) + .collect(); + tokio::try_join!( + async { + run_git( + None, + &["-c", AUTH, "clone", &url, path_str(&concurrent[0])?], + ) + .await + }, + async { + run_git( + None, + &["-c", AUTH, "clone", &url, path_str(&concurrent[1])?], + ) + .await + }, + async { + run_git( + None, + &["-c", AUTH, "clone", &url, path_str(&concurrent[2])?], + ) + .await + }, + async { + run_git( + None, + &["-c", AUTH, "clone", &url, path_str(&concurrent[3])?], + ) + .await + }, + )?; + for clone in &concurrent { + run_git(Some(clone), &["fsck", "--strict", "--full"]).await?; + assert_eq!(tokio::fs::read(clone.join("file")).await?, b"revision 3\n"); + } for protocol in ["0", "1", "2"] { let clone = workspace.path().join(format!("clone-{protocol}")); run_git( From ed6f27082c4316563039967323c5409b0fd3d9c6 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 15:52:22 -0700 Subject: [PATCH 02/12] Receive concurrent push uploads before serial publication --- docs/performance-plan.md | 6 +++ src/git_gateway.rs | 5 ++- tests/smart_http.rs | 3 ++ tests/smart_http/push_uploads.rs | 71 ++++++++++++++++++++++++++++++++ 4 files changed, 84 insertions(+), 1 deletion(-) create mode 100644 tests/smart_http/push_uploads.rs diff --git a/docs/performance-plan.md b/docs/performance-plan.md index 28a2e6b..f3661be 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -77,6 +77,12 @@ objects publish atomically; writers for the same object ID converge while different IDs may hydrate in parallel. A cold concurrent stock-Git clone test checks complete packs and `git fsck`. This removes one serialization point; it does not establish hot-repository throughput or a 10,000-repository result. +Push upload spooling uses private, budgeted scratch files and can now proceed +concurrently for one repository. The push transaction—from idempotency check +and gzip decoding through native Git execution and durable publication—remains +serialized per repository. A stalled upload regression test checks that a +second push can finish receiving before the first upload completes; it does +not establish parallel publication capacity. The native walk still traverses requested history on warm requests; native Git's own traversal memory is not bounded by the Rust batch size. Structural hydration pages certified edges and retains a visited set proportional to the requested diff --git a/src/git_gateway.rs b/src/git_gateway.rs index 4dd33f9..5c021e6 100644 --- a/src/git_gateway.rs +++ b/src/git_gateway.rs @@ -190,10 +190,13 @@ impl GitGateway { if !request.authenticated { return Err(GatewayError::Unauthorized); } - let _push = self.push.lock().await; let request = self.receive(request, None, admission).await?; let id = push_id.unwrap_or_else(|| uuid::Uuid::new_v4().into_bytes()); let digest = request_digest(&request).await?; + // Upload spooling uses a private, budgeted scratch file. Serialize + // the push-ID check, decode, native Git work and publication, but + // do not let one slow client block another client's upload. + let _push = self.push.lock().await; if self.repository.begin_push(id, actor, digest).await? { return Ok(http_body(with_push_id( self.repository.completed_response(id).await?, diff --git a/tests/smart_http.rs b/tests/smart_http.rs index 57efcc0..665414e 100644 --- a/tests/smart_http.rs +++ b/tests/smart_http.rs @@ -32,6 +32,8 @@ mod encoded_input; mod native_resources; #[path = "smart_http/publication.rs"] mod publication; +#[path = "smart_http/push_uploads.rs"] +mod push_uploads; #[path = "smart_http/ref_snapshots.rs"] mod ref_snapshots; @@ -178,6 +180,7 @@ async fn stock_git_push_and_clone_are_backed_by_one_repository_cell() Err(LfsError::Forbidden) )); assert!(repository.lfs_object(revoked_oid).await?.output.is_none()); + push_uploads::verify(&gateway).await?; let listener = TcpListener::bind("127.0.0.1:0").await?; let address = listener.local_addr()?; let teardown_gateway = Arc::clone(&gateway); diff --git a/tests/smart_http/push_uploads.rs b/tests/smart_http/push_uploads.rs new file mode 100644 index 0000000..a3cc36d --- /dev/null +++ b/tests/smart_http/push_uploads.rs @@ -0,0 +1,71 @@ +use super::*; +use canopy_server::git_http::GitHttpRequest; +use futures_core::Stream; +use std::{ + future::Future, + io, + pin::Pin, + task::{Context, Poll}, + time::Duration, +}; + +struct PausedUpload { + entered: Option>, + release: oneshot::Receiver<()>, +} + +impl Stream for PausedUpload { + type Item = Result; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + if let Some(entered) = self.entered.take() { + let _ = entered.send(()); + } + match Pin::new(&mut self.release).poll(cx) { + Poll::Pending => Poll::Pending, + Poll::Ready(_) => Poll::Ready(None), + } + } +} + +fn push(body: Body) -> GitHttpRequest { + GitHttpRequest { + method: "POST".into(), + path_info: "/repo.git/git-receive-pack".into(), + query: String::new(), + content_type: Some("application/x-git-receive-pack-request".into()), + gzip: false, + protocol_v2: false, + authenticated: true, + body, + } +} + +pub async fn verify(gateway: &Arc) -> Result<(), Box> { + let (entered_tx, entered_rx) = oneshot::channel(); + let (release_tx, release_rx) = oneshot::channel(); + let first_gateway = Arc::clone(gateway); + let stalled = tokio::spawn(async move { + first_gateway + .handle( + push(Body::from_stream(PausedUpload { + entered: Some(entered_tx), + release: release_rx, + })), + "canopy", + None, + None, + ) + .await + }); + tokio::time::timeout(Duration::from_secs(5), entered_rx).await??; + let second = tokio::time::timeout( + Duration::from_secs(5), + gateway.handle(push(Body::from("invalid")), "canopy", None, None), + ) + .await; + let _ = release_tx.send(()); + assert!(stalled.await?.is_err()); + assert!(second?.is_err()); + Ok(()) +} From c9a553cfdb97211ae2521806cf6312b3443af949 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 16:04:26 -0700 Subject: [PATCH 03/12] Benchmark stock Git transfers and unique-ref pushes across gateways --- docs/performance-plan.md | 54 +++++++++--- scripts/benchmark_repositories.py | 109 ++++++++++++++++++++----- scripts/test_benchmark_repositories.py | 77 +++++++++++++++++ 3 files changed, 210 insertions(+), 30 deletions(-) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index f3661be..97b919f 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -427,13 +427,42 @@ python3 -B scripts/benchmark_repositories.py \ run --active-repositories 1000 --distribution uniform \ --operation refs --rate 20 --duration 120 --concurrency 32 \ --output /path/to/canopy-two-gateway-refs.json +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-corpus.json \ + run --active-repositories 100 --distribution skewed \ + --operation clone --rate 2 --duration 120 --concurrency 8 \ + --work-dir /path/to/new-client-clone-scratch \ + --output /path/to/canopy-two-gateway-clones.json +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-corpus.json \ + run --active-repositories 1000 --distribution uniform \ + --operation push_branch --rate 5 --duration 120 --concurrency 8 \ + --work-dir /path/to/new-client-push-source \ + --output /path/to/canopy-two-gateway-pushes.json ``` The run assigns scheduled requests to ingresses by sequence number and records -each ingress's outcomes and latency percentiles. Compare the same workload with one and two ingresses -and collect node, object-store, and client resource metrics separately. This -driver measures metadata or Git v2 discovery; it does not measure push, pack -generation, full fetch, or LFS transfer throughput. +each ingress's outcomes and latency percentiles. Compare the same workload with +one and two ingresses and collect node, object-store, and client resource metrics +separately. `clone` measures a stock-Git full clone and verifies the seeded commit +and README hash. `cold_fetch` initializes an empty bare client and fetches +`refs/heads/main`; substitute it for `clone` with a separate fresh client work +directory and report. Both select only populated repositories, so the sample +above exercises 100 of the 10,000 identities. Disposable client checkouts are +removed after each attempt; the caller-supplied work directory itself remains. +Client Git and disk work is included in latency. `push_branch` creates one local +commit and pushes it to a unique `refs/heads/canopy-benchmark//` +ref for each arrival. The report records the run ID and commit; pushed refs +remain in the disposable corpus for inspection and must be included in its +eventual cleanup. The first push to a repository transfers objects, while later +pushes of the same commit mainly measure ref publication. Source commit setup is +outside the scheduled interval. Interpret those as +different workloads. These runs do not measure incremental fetch, pull, LFS, +or large-history throughput. Run `--operation refs` for Git v2 discovery or `--distribution skewed` for 90% of requests to the selected working set's first tenth. A fixed seed determines @@ -441,13 +470,14 @@ working-set selection and offered arrivals. The driver never calls a working set warm automatically: its count is not the server's resident count. Prewarm a set that fits the node before claiming warm latency, or label the run as cold/mixed. -Each worker reuses an HTTP connection. Arrivals follow a fixed clock schedule; +Each HTTP worker reuses a connection; stock-Git attempts use fresh client +repositories. Arrivals follow a fixed clock schedule; end-to-end latency starts at the scheduled instant and includes driver dispatch delay and complete response consumption. A bounded semaphore limits outstanding requests. Saturation records `driver_busy` instead of accumulating an unbounded -client queue. HTTP errors and transport failures are recorded without retries. -`--timeout` bounds individual socket operations; it is not a whole-request -deadline. All outcomes go to a sibling `.samples.jsonl`; the JSON summary includes counts, +client queue. HTTP and Git failures are recorded without retries. `--timeout` +bounds individual HTTP socket operations; `--git-timeout` bounds each Git process. +All outcomes go to a sibling `.samples.jsonl`; the JSON summary includes counts, error totals, scheduled/service/dispatch latency percentiles and the manifest SHA-256. `started_at_utc` anchors the run to server logs; each sample's scheduled offset is `sequence / offered_rps`, while latency continues to use the @@ -455,14 +485,16 @@ monotonic clock. Dropped arrivals have no fabricated zero latency. Any failure y exit status 1 after reports are written. Latency percentiles include completed errors; always read them alongside error/drop counts. -The driver unit test runs a real HTTP server that deliberately rejects requests: +The driver tests run a real HTTP server that deliberately rejects requests and +stock Git against a disposable local bare repository: ```sh python3 -B -m unittest discover -s scripts -p test_benchmark_repositories.py -v ``` -It verifies concurrency bounds, complete scheduled-outcome accounting, absence -of retries, queue-delay inclusion and credential exclusion from the report. +They verify concurrency bounds, complete scheduled-outcome accounting, absence +of retries, queue-delay inclusion, Git tip/content checks and credential +exclusion from the report. ## Measurement history diff --git a/scripts/benchmark_repositories.py b/scripts/benchmark_repositories.py index 7e9d4a4..52751df 100644 --- a/scripts/benchmark_repositories.py +++ b/scripts/benchmark_repositories.py @@ -1,9 +1,10 @@ #!/usr/bin/env python3 -"""Seed disposable repository identities and measure scheduled HTTP read load. +"""Seed disposable repositories and measure scheduled HTTP and stock-Git work. Use CANOPY_GIT_TOKEN for authentication. Reports contain no credentials. This -measures metadata or Git v2 discovery; clone, push and LFS throughput need their -own qualification. Seed and verify use stock Git for a declared corpus sample. +measures metadata, Git v2 discovery, full clone, cold fetch or unique-ref push; +incremental fetch, pull and LFS throughput need their own qualification. Seed +and verify use stock Git for a declared corpus sample. """ import argparse @@ -19,6 +20,7 @@ import random import re import subprocess +import tempfile import threading import time from urllib.parse import quote, urlsplit @@ -31,6 +33,7 @@ def __init__(self, base, token, timeout): if url.scheme not in ("http", "https") or not url.hostname or url.username or url.password or url.query or url.fragment: raise ValueError("base URL must be HTTP(S), without credentials, query or fragment") self.connection_type = http.client.HTTPSConnection if url.scheme == "https" else http.client.HTTPConnection + self.base_url = base.rstrip("/") self.host, self.port, self.prefix = url.hostname, url.port, url.path.rstrip("/") self.token, self.timeout = token, timeout self.local = threading.local() @@ -67,7 +70,7 @@ def close(self): connection.close() -def git(*args, cwd, token): +def git(*args, cwd, token, timeout=120, request_id=None): environment = {key: value for key, value in os.environ.items() if not key.startswith(("GIT_", "AWS_", "RUSTFS_", "CANOPY_"))} environment.update(GIT_TERMINAL_PROMPT="0", GIT_CONFIG_NOSYSTEM="1", @@ -75,8 +78,11 @@ def git(*args, cwd, token): GIT_CONFIG_KEY_0="credential.helper", GIT_CONFIG_VALUE_0="", GIT_CONFIG_KEY_1="http.extraHeader", GIT_CONFIG_VALUE_1=f"Authorization: Bearer {token}") + if request_id is not None: + environment.update(GIT_CONFIG_COUNT="3", GIT_CONFIG_KEY_2="http.extraHeader", + GIT_CONFIG_VALUE_2=f"X-Request-ID: {request_id}") result = subprocess.run(["git", *args], cwd=cwd, env=environment, - capture_output=True, timeout=120, check=False) + capture_output=True, timeout=timeout, check=False) if result.returncode: raise RuntimeError(f"Git {args[0]} failed (exit {result.returncode})") return result.stdout.strip().decode() @@ -182,19 +188,68 @@ def percentiles(values): for name, q in (("p50", .5), ("p95", .95), ("p99", .99), ("max", 1))} +def git_transfer(operation, url, entry, token, work_dir, timeout, request_id): + """Run one disposable stock-Git transfer and validate its advertised tip.""" + with tempfile.TemporaryDirectory(prefix="canopy-git-read-", dir=work_dir) as temporary: + destination = Path(temporary) / "repo" + if operation == "clone": + git("clone", "--quiet", url, str(destination), cwd=temporary, token=token, + timeout=timeout, request_id=request_id) + tip = git("rev-parse", "HEAD", cwd=destination, token=token, timeout=timeout) + readme = destination / "README.md" + return tip == entry["commit"] and readme.is_file() and ( + hashlib.sha256(readme.read_bytes()).hexdigest() == entry["readme_sha256"]) + if operation == "cold_fetch": + git("init", "--bare", "--quiet", str(destination), cwd=temporary, + token=token, timeout=timeout) + git("fetch", "--quiet", url, "refs/heads/main", cwd=destination, + token=token, timeout=timeout, request_id=request_id) + return git("rev-parse", "FETCH_HEAD", cwd=destination, token=token, + timeout=timeout) == entry["commit"] + raise ValueError("unsupported Git transfer operation") + + +def push_branch(url, source, reference, token, timeout, request_id): + """Publish one unique benchmark ref; Git checks the receive-pack report.""" + git("push", "--quiet", url, f"HEAD:{reference}", cwd=source, token=token, + timeout=timeout, request_id=request_id) + + def measure(args, client, token): - del token clients = client if isinstance(client, list) else [client] manifest = corpus(args.manifest) - if args.active_repositories > len(manifest["repositories"]): - raise ValueError("active repository count exceeds corpus") - generator = random.Random(args.seed) - active = generator.sample(manifest["repositories"], args.active_repositories) - # Skew has a declared hot tenth, not an implicit warm-cache assumption. - hot = active[:max(1, len(active) // 10)] + git_read = args.operation in ("clone", "cold_fetch") + git_write = args.operation == "push_branch" + git_operation = git_read or git_write + eligible = ([entry for entry in manifest["repositories"] if entry["commit"] is not None] + if git_read else manifest["repositories"]) + if args.active_repositories > len(eligible): + raise ValueError("active repository count exceeds eligible corpus") total = math.ceil(args.duration * args.rate) if total > 1_000_000: raise ValueError("one run is limited to one million scheduled arrivals") + samples_path = args.output.with_suffix(".samples.jsonl") + if args.output.exists() or samples_path.exists(): + raise ValueError("run requires new output paths") + if git_operation: + if args.work_dir is None: + raise ValueError("stock-Git runs require --work-dir") + args.work_dir.mkdir(parents=True, exist_ok=False) + push_run_id = uuid.uuid4().hex if git_write else None + source = args.work_dir / "source" if git_write else None + push_commit = None + if git_write: + git("init", "-b", "main", str(source), cwd=args.work_dir, token=token) + git("config", "user.name", "Canopy Benchmark", cwd=source, token=token) + git("config", "user.email", "benchmark@example.invalid", cwd=source, token=token) + (source / "README.md").write_text(f"benchmark run {push_run_id}\n") + git("add", "README.md", cwd=source, token=token) + git("commit", "-m", "Benchmark fixture", cwd=source, token=token) + push_commit = git("rev-parse", "HEAD", cwd=source, token=token) + generator = random.Random(args.seed) + active = generator.sample(eligible, args.active_repositories) + # Skew has a declared hot tenth, not an implicit warm-cache assumption. + hot = active[:max(1, len(active) // 10)] counts, latencies, service_times, dispatch_times = Counter(), [], [], [] ingress_counts = [Counter() for _ in clients] ingress_latencies = [[] for _ in clients] @@ -203,9 +258,6 @@ def measure(args, client, token): lock = threading.Lock() started = time.monotonic() started_at = datetime.now(timezone.utc).isoformat() - samples_path = args.output.with_suffix(".samples.jsonl") - if args.output.exists() or samples_path.exists(): - raise ValueError("run requires new output paths") with samples_path.open("x") as samples: def record(sample): with lock: @@ -230,10 +282,24 @@ def execute(sequence, entry, scheduled): status, body = ingress.request(f"/api/repositories/{entry['name']}", request_id=request_id) decoded = json.loads(body) if status == 200 else None valid = isinstance(decoded, dict) and decoded.get("repository_id") == entry["repository_id"] - else: + elif args.operation == "refs": status, body = ingress.request(f"/{entry['owner']}/{entry['name']}.git/info/refs?service=git-upload-pack", git=True, request_id=request_id) valid = status == 200 and body.startswith(b"000eversion 2\n") - result = "ok" if valid else (f"http_{status}" if status != 200 else "invalid_response") + elif git_read: + url = f"{ingress.base_url}/{entry['owner']}/{entry['name']}.git" + valid = git_transfer(args.operation, url, entry, token, args.work_dir, + args.git_timeout, request_id) + else: + url = f"{ingress.base_url}/{entry['owner']}/{entry['name']}.git" + reference = f"refs/heads/canopy-benchmark/{push_run_id}/{sequence:07d}" + push_branch(url, source, reference, token, args.git_timeout, request_id) + valid = True + result = ("ok" if valid else + (f"http_{status}" if not git_operation and status != 200 else "invalid_response")) + except subprocess.TimeoutExpired: + result = "client_timeout" + except RuntimeError: + result = "git_error" except (OSError, ValueError, KeyError, http.client.HTTPException): pass finally: @@ -261,11 +327,14 @@ def execute(sequence, entry, scheduled): elapsed = time.monotonic() - started result = {"version": 1, "started_at_utc": started_at, "corpus_repositories": len(manifest["repositories"]), + "eligible_repositories": len(eligible), "active_repositories": len(active), "distribution": args.distribution, "operation": args.operation, "seed": args.seed, "offered_rps": args.rate, "schedule_seconds": args.duration, "elapsed_including_drain_seconds": round(elapsed, 3), "concurrency": args.concurrency, "request_timeout_seconds": args.timeout, + "git_timeout_seconds": args.git_timeout if git_operation else None, + "push_run_id": push_run_id, "push_commit": push_commit, "scheduled": total, "outcomes": dict(counts), "ingresses": [{"index": index, "outcomes": dict(outcomes), "scheduled_latency_ms": percentiles(ingress_latencies[index]), @@ -274,7 +343,7 @@ def execute(sequence, entry, scheduled): "failed_arrivals": total - counts["ok"], "scheduled_latency_ms": percentiles(latencies), "service_ms": percentiles(service_times), "dispatch_delay_ms": percentiles(dispatch_times), - "latency_population": "all completed HTTP attempts, including errors; driver_busy arrivals are counted failures without fabricated latency", + "latency_population": "all completed HTTP or stock-Git attempts, including errors and client validation; driver_busy arrivals are counted failures without fabricated latency", "manifest_sha256": hashlib.sha256(args.manifest.read_bytes()).hexdigest()} if sum(counts.values()) != total: raise RuntimeError("benchmark lost scheduled outcomes") @@ -307,10 +376,12 @@ def main(): run = commands.add_parser("run") run.add_argument("--active-repositories", type=positive, required=True) run.add_argument("--distribution", choices=("uniform", "skewed"), default="uniform") - run.add_argument("--operation", choices=("metadata", "refs"), default="metadata") + run.add_argument("--operation", choices=("metadata", "refs", "clone", "cold_fetch", "push_branch"), default="metadata") run.add_argument("--rate", type=positive, default=20) run.add_argument("--duration", type=positive, default=30) run.add_argument("--concurrency", type=positive, default=32) + run.add_argument("--work-dir", type=Path, help="new client directory for stock-Git operations") + run.add_argument("--git-timeout", type=positive, default=120) run.add_argument("--output", type=Path, required=True) args = parser.parse_args() token = os.environ.get("CANOPY_GIT_TOKEN") diff --git a/scripts/test_benchmark_repositories.py b/scripts/test_benchmark_repositories.py index 9c3a0f0..201be72 100644 --- a/scripts/test_benchmark_repositories.py +++ b/scripts/test_benchmark_repositories.py @@ -1,6 +1,7 @@ """Verify the capacity driver against an overloaded real HTTP fixture.""" from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +import hashlib import json from pathlib import Path import tempfile @@ -23,6 +24,7 @@ def do_GET(self): with self.server.guard: self.server.requests += 1 self.server.request_ids.append(self.headers.get("X-Request-ID")) + self.server.auth_headers.append(self.headers.get("Authorization")) self.server.active += 1 self.server.peak = max(self.server.peak, self.server.active) try: @@ -43,6 +45,7 @@ def test_overload_has_no_hidden_retries_or_missing_arrivals(self): server.guard = threading.Lock() server.requests = server.active = server.peak = 0 server.request_ids = [] + server.auth_headers = [] thread = threading.Thread(target=server.serve_forever, daemon=True) thread.start() client = benchmark.Client(f"http://127.0.0.1:{server.server_port}", "fixture-token", 2) @@ -86,6 +89,7 @@ def test_multiple_ingresses_report_each_gateway_without_losing_arrivals(self): server.guard = threading.Lock() server.requests = server.active = server.peak = 0 server.request_ids = [] + server.auth_headers = [] thread = threading.Thread(target=server.serve_forever, daemon=True) thread.start() threads.append(thread) @@ -122,6 +126,79 @@ def test_multiple_ingresses_report_each_gateway_without_losing_arrivals(self): for thread in threads: thread.join() + def test_stock_git_sends_authorization_and_request_id_headers(self): + server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) + server.guard = threading.Lock() + server.requests = server.active = server.peak = 0 + server.request_ids = [] + server.auth_headers = [] + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + try: + with tempfile.TemporaryDirectory() as directory: + request_id = str(uuid.uuid4()) + with self.assertRaises(RuntimeError): + benchmark.git("ls-remote", f"http://127.0.0.1:{server.server_port}/repo.git", + cwd=directory, token="fixture-token", request_id=request_id) + self.assertGreater(server.requests, 0) + self.assertEqual(server.request_ids, [request_id] * server.requests) + self.assertEqual(server.auth_headers, ["Bearer fixture-token"] * server.requests) + finally: + server.shutdown() + server.server_close() + thread.join() + + def test_stock_git_clone_and_cold_fetch_validate_populated_repositories(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + source = root / "source" + remote = root / "canopy" / "fixture.git" + remote.parent.mkdir() + benchmark.git("init", "-b", "main", str(source), cwd=root, token="fixture-token") + benchmark.git("config", "user.name", "Fixture", cwd=source, token="fixture-token") + benchmark.git("config", "user.email", "fixture@example.invalid", cwd=source, + token="fixture-token") + body = b"stock Git transfer fixture\n" + (source / "README.md").write_bytes(body) + benchmark.git("add", "README.md", cwd=source, token="fixture-token") + benchmark.git("commit", "-m", "Fixture", cwd=source, token="fixture-token") + benchmark.git("init", "--bare", "-b", "main", str(remote), cwd=root, + token="fixture-token") + benchmark.git("push", str(remote), "HEAD:refs/heads/main", cwd=source, + token="fixture-token") + entry = {"name": "fixture", "owner": "canopy", + "repository_id": str(uuid.uuid4()), + "commit": benchmark.git("rev-parse", "HEAD", cwd=source, + token="fixture-token"), + "readme_sha256": hashlib.sha256(body).hexdigest()} + manifest = root / "manifest.json" + manifest.write_text(json.dumps({"version": 1, "complete": True, + "requested_repositories": 1, "repositories": [entry]})) + ingresses = [SimpleNamespace(base_url=root.as_uri()) for _ in range(2)] + for operation in ("clone", "cold_fetch"): + args = SimpleNamespace(manifest=manifest, active_repositories=1, seed=42, + duration=1, rate=2, concurrency=2, distribution="uniform", + operation=operation, timeout=2, git_timeout=30, + work_dir=root / f"{operation}-scratch", output=root / f"{operation}.json") + report = benchmark.measure(args, ingresses, "fixture-token") + self.assertEqual(report["outcomes"], {"ok": 2}) + self.assertEqual(report["eligible_repositories"], 1) + self.assertEqual([item["outcomes"] for item in report["ingresses"]], + [{"ok": 1}, {"ok": 1}]) + self.assertEqual(list(args.work_dir.iterdir()), []) + self.assertNotIn("fixture-token", args.output.read_text()) + args = SimpleNamespace(manifest=manifest, active_repositories=1, seed=42, + duration=1, rate=2, concurrency=2, distribution="uniform", + operation="push_branch", timeout=2, git_timeout=30, + work_dir=root / "push-scratch", output=root / "push.json") + report = benchmark.measure(args, ingresses, "fixture-token") + self.assertEqual(report["outcomes"], {"ok": 2}) + refs = benchmark.git("for-each-ref", "--format=%(objectname)", + f"refs/heads/canopy-benchmark/{report['push_run_id']}", cwd=remote, + token="fixture-token").splitlines() + self.assertEqual(refs, [report["push_commit"]] * 2) + self.assertNotIn("fixture-token", args.output.read_text()) + if __name__ == "__main__": unittest.main() From 15d88c763ea2530864f63dcd33496153eb798420 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 16:06:53 -0700 Subject: [PATCH 04/12] Allocate Git object write coordination only for hydrating caches --- docs/performance-plan.md | 4 +++- src/git_cache.rs | 22 ++++++++++++---------- src/git_cache/tests.rs | 3 +++ 3 files changed, 18 insertions(+), 11 deletions(-) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index 97b919f..a01464b 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -74,7 +74,9 @@ Preparation logs report newly hydrated blob counts/bytes and elapsed time. Concurrent fetches of one repository now share the immutable object cache without holding its gateway mutex for the entire preparation. Verified loose objects publish atomically; writers for the same object ID converge while -different IDs may hydrate in parallel. A cold concurrent stock-Git clone test +different IDs may hydrate in parallel. Writer locks allocate only when an object +cache first hydrates an object, not for each disposable ref snapshot. A cold +concurrent stock-Git clone test checks complete packs and `git fsck`. This removes one serialization point; it does not establish hot-repository throughput or a 10,000-repository result. Push upload spooling uses private, budgeted scratch files and can now proceed diff --git a/src/git_cache.rs b/src/git_cache.rs index 997485d..1a0ba85 100644 --- a/src/git_cache.rs +++ b/src/git_cache.rs @@ -5,7 +5,7 @@ use std::{ fs::{self, File, OpenOptions}, io::{self, BufWriter, Write}, path::{Path, PathBuf}, - sync::Arc, + sync::{Arc, OnceLock}, }; use cellule_ltx::{DiskBudget, DiskReservation, LtxError}; @@ -54,7 +54,7 @@ pub(crate) struct GitCache { objects: Option>, // Only durable hydration writes this cache. Stripe by OID so concurrent // fetches share a completed loose object without serializing all objects. - object_writes: [Arc>; 64], + object_writes: OnceLock<[Arc>; 64]>, } impl GitCache { @@ -86,7 +86,7 @@ impl GitCache { directory: tempfile::Builder::new().prefix(CACHE_PREFIX).tempdir_in(fs::canonicalize(root)?)?, reservation: Some(budget.try_reserve(0)?), objects, - object_writes: std::array::from_fn(|_| Arc::new(Mutex::new(()))), + object_writes: OnceLock::new(), }); for directory in ["objects/info", "objects/pack", "refs/heads", "refs/tags", "hooks"] { fs::create_dir_all(cache.git_dir().join(directory))?; @@ -181,6 +181,13 @@ impl GitCache { } } + fn object_write_lock(&self, oid: crate::ObjectId) -> Arc> { + let stripes = self + .object_writes + .get_or_init(|| std::array::from_fn(|_| Arc::new(Mutex::new(())))); + Arc::clone(&stripes[oid[0] as usize % stripes.len()]) + } + fn object_writer( self: &Arc, oid: crate::ObjectId, @@ -245,9 +252,7 @@ impl GitCache { kind: ObjectKind, body: Vec, ) -> Result<(), CacheError> { - let write = Arc::clone(&self.object_writes[oid[0] as usize % self.object_writes.len()]) - .lock_owned() - .await; + let write = self.object_write_lock(oid).lock_owned().await; let cache = Arc::clone(self); tokio::task::spawn_blocking(move || { let _write = write; @@ -274,10 +279,7 @@ impl GitCache { mut reader: LargeBlobRead, ) -> Result<(), CacheError> { let reference = reader.reference(); - let write = - Arc::clone(&self.object_writes[reference.oid[0] as usize % self.object_writes.len()]) - .lock_owned() - .await; + let write = self.object_write_lock(reference.oid).lock_owned().await; let cache = Arc::clone(self); let pending = tokio::task::spawn_blocking(move || { if reference.oid.format() != cache.object_format { diff --git a/src/git_cache/tests.rs b/src/git_cache/tests.rs index cc77051..09c775c 100644 --- a/src/git_cache/tests.rs +++ b/src/git_cache/tests.rs @@ -12,6 +12,7 @@ async fn concurrent_hydration_publishes_each_object_once() -> Result<(), Box Result<(), Box Date: Tue, 29 Sep 2026 16:15:44 -0700 Subject: [PATCH 05/12] Benchmark incremental fetch and pull with versioned corpus --- docs/performance-plan.md | 61 ++++++++++++-- scripts/benchmark_repositories.py | 110 ++++++++++++++++++++++--- scripts/test_benchmark_repositories.py | 75 +++++++++++++++++ 3 files changed, 225 insertions(+), 21 deletions(-) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index a01464b..58a0287 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -254,10 +254,12 @@ Acceptance: paused ownership lookup and release do not block an unrelated warm repository's metadata or stock Git discovery. Concurrent requests for the same cold repository converge on one serving entry. All pins, denied/lost release, cleanup failure, cancellation and shutdown tests continue to pass. The lock -isolation part is implemented. The initial density driver below measures metadata -and Git v2 discovery, with stock Git sampling and full identity recovery checks. -Transition queue/service timings are logged. Other workload classes, comprehensive -resource metrics and independent load-generator deployment remain open. +isolation part is implemented. The automated bounded-container run still measures +metadata and Git v2 discovery, with stock Git sampling and full identity recovery +checks. The standalone driver can schedule clone, cold/incremental fetch, pull and +unique-ref push, but those workloads have not been qualified at target scale. +Transition queue/service timings are logged. Comprehensive resource metrics and +independent load-generator deployment remain open. ### 2. Measured active-Cell admission @@ -412,6 +414,10 @@ The default sample contains three one-commit repositories; the remainder are empty. This is a density smoke corpus, not realistic large-history qualification. An interrupted seed leaves an explicitly incomplete manifest and cannot be used as a complete corpus. Names have a unique prefix; the manifest records them. +The default seed remains a version-1, one-commit corpus. Use +`seed --incremental-fixture` with a new manifest and storage prefix to create +version 2: each populated repository has a `benchmark-base` ref at its first +commit and `main` one commit ahead. `verify` accepts both versions. For a multi-gateway read run, seed the documented 10,000-repository corpus, start two nodes against the same deployment, and pass each ingress to `run`: @@ -462,9 +468,46 @@ ref for each arrival. The report records the run ID and commit; pushed refs remain in the disposable corpus for inspection and must be included in its eventual cleanup. The first push to a repository transfers objects, while later pushes of the same commit mainly measure ref publication. Source commit setup is -outside the scheduled interval. Interpret those as -different workloads. These runs do not measure incremental fetch, pull, LFS, -or large-history throughput. +outside the scheduled interval. Interpret those as different workloads. The +version-1 runs above do not measure incremental fetch, pull, LFS, or +large-history throughput. + +For an incremental fetch and pull run, use a separate version-2 corpus. The +driver first fetches each selected base ref into a local bare template, outside +the arrival clock, then gives every arrival a fresh client that shares that +base object store. `incremental_fetch` fetches `main` into a bare client; +`incremental_pull` fast-forwards a checked-out client and verifies the new +file. Use distinct fresh work directories and output paths: + +```bash +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --manifest /path/to/canopy-incremental-corpus.json \ + seed --repositories 10000 --populated 100 --incremental-fixture \ + --work-dir /path/to/canopy-incremental-seed +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-incremental-corpus.json \ + run --active-repositories 100 --operation incremental_fetch \ + --rate 2 --duration 120 --concurrency 8 \ + --work-dir /path/to/new-incremental-fetch-clients \ + --output /path/to/canopy-incremental-fetch.json +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-incremental-corpus.json \ + run --active-repositories 100 --operation incremental_pull \ + --rate 2 --duration 120 --concurrency 8 \ + --work-dir /path/to/new-incremental-pull-clients \ + --output /path/to/canopy-incremental-pull.json +``` + +Template preparation is real server traffic and can warm both gateways; its +elapsed time is recorded separately. These runs measure one-commit incremental +transfers, not large histories or cold-cache startup. Client-side shared clone +setup is included in each attempt's latency. Report node and generator resource +usage alongside the latency distribution. Run `--operation refs` for Git v2 discovery or `--distribution skewed` for 90% of requests to the selected working set's first tenth. A fixed seed determines @@ -495,8 +538,8 @@ python3 -B -m unittest discover -s scripts -p test_benchmark_repositories.py -v ``` They verify concurrency bounds, complete scheduled-outcome accounting, absence -of retries, queue-delay inclusion, Git tip/content checks and credential -exclusion from the report. +of retries, queue-delay inclusion, version-1/version-2 corpus and Git +tip/content checks, and credential exclusion from the report. ## Measurement history diff --git a/scripts/benchmark_repositories.py b/scripts/benchmark_repositories.py index 52751df..58ae959 100644 --- a/scripts/benchmark_repositories.py +++ b/scripts/benchmark_repositories.py @@ -2,9 +2,10 @@ """Seed disposable repositories and measure scheduled HTTP and stock-Git work. Use CANOPY_GIT_TOKEN for authentication. Reports contain no credentials. This -measures metadata, Git v2 discovery, full clone, cold fetch or unique-ref push; -incremental fetch, pull and LFS throughput need their own qualification. Seed -and verify use stock Git for a declared corpus sample. +measures metadata, Git v2 discovery, clone, fetch, pull or unique-ref push; +LFS throughput needs its own qualification. Seed and verify use stock Git for +a declared corpus sample. Incremental workloads require an opt-in two-commit +corpus. """ import argparse @@ -101,7 +102,8 @@ def seed(args, client, token): generator = random.Random(args.seed) selected = set(generator.sample(range(args.repositories), min(args.populated, args.repositories))) prefix = f"density-{uuid.uuid4().hex[:12]}" - manifest = {"version": 1, "complete": False, "prefix": prefix, "seed": args.seed, + manifest = {"version": 2 if args.incremental_fixture else 1, + "complete": False, "prefix": prefix, "seed": args.seed, "requested_repositories": args.repositories, "repositories": [], "git_version": git("--version", cwd=args.work_dir, token=token)} started = time.monotonic() @@ -120,6 +122,8 @@ def seed(args, client, token): record = {key: entry[key] for key in ("name", "owner", "repository_id")} record["create_ms"] = round(create_ms, 3) record["commit"] = None + if args.incremental_fixture: + record["base_commit"] = None manifest["repositories"].append(record) if index in selected: local = args.work_dir / name @@ -130,8 +134,18 @@ def seed(args, client, token): (local / "README.md").write_bytes(content) git("add", "README.md", cwd=local, token=token) git("commit", "-m", "Density fixture", cwd=local, token=token) + if args.incremental_fixture: + record["base_commit"] = git("rev-parse", "HEAD", cwd=local, token=token) + git("branch", "benchmark-base", cwd=local, token=token) + increment = f"{name}\nincremental={args.seed}\n".encode() + (local / "incremental.txt").write_bytes(increment) + git("add", "incremental.txt", cwd=local, token=token) + git("commit", "-m", "Incremental fixture", cwd=local, token=token) + record["incremental_sha256"] = hashlib.sha256(increment).hexdigest() url = f"{args.base_url.rstrip('/')}/{entry['owner']}/{name}.git" - git("push", url, "HEAD:refs/heads/main", cwd=local, token=token) + refs = (["HEAD:refs/heads/main", "benchmark-base:refs/heads/benchmark-base"] + if args.incremental_fixture else ["HEAD:refs/heads/main"]) + git("push", url, *refs, cwd=local, token=token) record["commit"] = git("rev-parse", "HEAD", cwd=local, token=token) record["readme_sha256"] = hashlib.sha256(content).hexdigest() if (index + 1) % 25 == 0: @@ -148,12 +162,19 @@ def seed(args, client, token): def corpus(path): manifest = json.loads(path.read_text()) entries = manifest["repositories"] - if manifest["version"] != 1 or not manifest["complete"] or len(entries) != manifest["requested_repositories"]: - raise ValueError("manifest must contain a complete version-1 corpus") + if manifest["version"] not in (1, 2) or not manifest["complete"] or len(entries) != manifest["requested_repositories"]: + raise ValueError("manifest must contain a complete version-1 or version-2 corpus") for entry in entries: if any(not re.fullmatch(r"[A-Za-z0-9_-]+", entry[key]) for key in ("name", "owner")): raise ValueError("manifest contains invalid repository names") uuid.UUID(entry["repository_id"]) + if manifest["version"] == 2 and entry["commit"] is not None: + base_commit = entry.get("base_commit") + digest = entry.get("incremental_sha256") + if not isinstance(base_commit, str) or not re.fullmatch(r"[0-9a-f]{40}|[0-9a-f]{64}", base_commit): + raise ValueError("version-2 corpus has an invalid base commit") + if not isinstance(digest, str) or not re.fullmatch(r"[0-9a-f]{64}", digest): + raise ValueError("version-2 corpus has an invalid incremental digest") return manifest @@ -177,6 +198,12 @@ def verify(args, client, token): raise RuntimeError("restored commit differs from manifest") if hashlib.sha256((clone / "README.md").read_bytes()).hexdigest() != entry["readme_sha256"]: raise RuntimeError("restored file bytes differ from manifest") + if manifest["version"] == 2: + if git("rev-parse", "refs/remotes/origin/benchmark-base", cwd=clone, + token=token) != entry["base_commit"]: + raise RuntimeError("restored incremental base differs from manifest") + if hashlib.sha256((clone / "incremental.txt").read_bytes()).hexdigest() != entry["incremental_sha256"]: + raise RuntimeError("restored incremental bytes differ from manifest") git("fsck", "--strict", "--full", cwd=clone, token=token) populated += 1 return {"verified_repositories": len(manifest["repositories"]), "git_v0_v2_samples": populated} @@ -188,7 +215,31 @@ def percentiles(values): for name, q in (("p50", .5), ("p95", .95), ("p99", .99), ("max", 1))} -def git_transfer(operation, url, entry, token, work_dir, timeout, request_id): +def prepare_incremental(active, clients, token, work_dir, timeout): + """Stage base-only Git repositories before the arrival clock starts.""" + root = work_dir / "incremental-templates" + root.mkdir() + templates = {} + for index, entry in enumerate(active): + destination = root / entry["repository_id"] + ingress = clients[index % len(clients)] + url = f"{ingress.base_url}/{entry['owner']}/{entry['name']}.git" + git("init", "--bare", "-b", "benchmark-base", str(destination), + cwd=root, token=token, timeout=timeout) + git("fetch", "--quiet", url, + "refs/heads/benchmark-base:refs/heads/benchmark-base", + cwd=destination, token=token, timeout=timeout) + if git("rev-parse", "HEAD", cwd=destination, token=token, + timeout=timeout) != entry["base_commit"]: + raise RuntimeError("incremental template differs from manifest") + templates[entry["repository_id"]] = destination + if (index + 1) % 25 == 0: + print(f"prepared {index + 1}/{len(active)} incremental clients", flush=True) + return templates + + +def git_transfer(operation, url, entry, token, work_dir, timeout, request_id, + template=None): """Run one disposable stock-Git transfer and validate its advertised tip.""" with tempfile.TemporaryDirectory(prefix="canopy-git-read-", dir=work_dir) as temporary: destination = Path(temporary) / "repo" @@ -206,6 +257,26 @@ def git_transfer(operation, url, entry, token, work_dir, timeout, request_id): token=token, timeout=timeout, request_id=request_id) return git("rev-parse", "FETCH_HEAD", cwd=destination, token=token, timeout=timeout) == entry["commit"] + if operation in ("incremental_fetch", "incremental_pull"): + if template is None: + raise ValueError("incremental transfer requires a prepared base") + clone_args = (["clone", "--quiet", "--shared", "--bare"] + if operation == "incremental_fetch" else + ["clone", "--quiet", "--shared"]) + git(*clone_args, str(template), str(destination), cwd=temporary, + token=token, timeout=timeout) + if operation == "incremental_fetch": + git("fetch", "--quiet", url, "refs/heads/main", cwd=destination, + token=token, timeout=timeout, request_id=request_id) + return git("rev-parse", "FETCH_HEAD", cwd=destination, + token=token, timeout=timeout) == entry["commit"] + git("pull", "--quiet", "--ff-only", url, "refs/heads/main", + cwd=destination, token=token, timeout=timeout, request_id=request_id) + incremental = destination / "incremental.txt" + return (git("rev-parse", "HEAD", cwd=destination, token=token, + timeout=timeout) == entry["commit"] and + incremental.is_file() and + hashlib.sha256(incremental.read_bytes()).hexdigest() == entry["incremental_sha256"]) raise ValueError("unsupported Git transfer operation") @@ -218,10 +289,13 @@ def push_branch(url, source, reference, token, timeout, request_id): def measure(args, client, token): clients = client if isinstance(client, list) else [client] manifest = corpus(args.manifest) - git_read = args.operation in ("clone", "cold_fetch") + incremental = args.operation in ("incremental_fetch", "incremental_pull") + git_read = args.operation in ("clone", "cold_fetch") or incremental git_write = args.operation == "push_branch" git_operation = git_read or git_write - eligible = ([entry for entry in manifest["repositories"] if entry["commit"] is not None] + eligible = ([entry for entry in manifest["repositories"] if entry.get("base_commit") is not None] + if incremental else + [entry for entry in manifest["repositories"] if entry["commit"] is not None] if git_read else manifest["repositories"]) if args.active_repositories > len(eligible): raise ValueError("active repository count exceeds eligible corpus") @@ -248,6 +322,12 @@ def measure(args, client, token): push_commit = git("rev-parse", "HEAD", cwd=source, token=token) generator = random.Random(args.seed) active = generator.sample(eligible, args.active_repositories) + templates = {} + setup_started = time.monotonic() + if incremental: + templates = prepare_incremental(active, clients, token, args.work_dir, + args.git_timeout) + setup_seconds = round(time.monotonic() - setup_started, 3) # Skew has a declared hot tenth, not an implicit warm-cache assumption. hot = active[:max(1, len(active) // 10)] counts, latencies, service_times, dispatch_times = Counter(), [], [], [] @@ -288,7 +368,8 @@ def execute(sequence, entry, scheduled): elif git_read: url = f"{ingress.base_url}/{entry['owner']}/{entry['name']}.git" valid = git_transfer(args.operation, url, entry, token, args.work_dir, - args.git_timeout, request_id) + args.git_timeout, request_id, + templates.get(entry["repository_id"])) else: url = f"{ingress.base_url}/{entry['owner']}/{entry['name']}.git" reference = f"refs/heads/canopy-benchmark/{push_run_id}/{sequence:07d}" @@ -332,6 +413,7 @@ def execute(sequence, entry, scheduled): "operation": args.operation, "seed": args.seed, "offered_rps": args.rate, "schedule_seconds": args.duration, "elapsed_including_drain_seconds": round(elapsed, 3), + "incremental_client_setup_seconds": setup_seconds if incremental else None, "concurrency": args.concurrency, "request_timeout_seconds": args.timeout, "git_timeout_seconds": args.git_timeout if git_operation else None, "push_run_id": push_run_id, "push_commit": push_commit, @@ -370,13 +452,17 @@ def main(): create = commands.add_parser("seed") create.add_argument("--repositories", type=positive, default=1000) create.add_argument("--populated", type=positive, default=3) + create.add_argument("--incremental-fixture", action="store_true", + help="seed a second commit and benchmark-base ref for incremental fetch/pull") create.add_argument("--work-dir", type=Path, required=True) check = commands.add_parser("verify") check.add_argument("--work-dir", type=Path, required=True) run = commands.add_parser("run") run.add_argument("--active-repositories", type=positive, required=True) run.add_argument("--distribution", choices=("uniform", "skewed"), default="uniform") - run.add_argument("--operation", choices=("metadata", "refs", "clone", "cold_fetch", "push_branch"), default="metadata") + run.add_argument("--operation", choices=("metadata", "refs", "clone", "cold_fetch", + "incremental_fetch", "incremental_pull", "push_branch"), + default="metadata") run.add_argument("--rate", type=positive, default=20) run.add_argument("--duration", type=positive, default=30) run.add_argument("--concurrency", type=positive, default=32) diff --git a/scripts/test_benchmark_repositories.py b/scripts/test_benchmark_repositories.py index 201be72..f80af18 100644 --- a/scripts/test_benchmark_repositories.py +++ b/scripts/test_benchmark_repositories.py @@ -40,6 +40,56 @@ def do_GET(self): class ScheduledLoad(unittest.TestCase): + def test_incremental_seed_and_verify_preserve_both_commits(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + remotes = root / "remotes" + remotes.mkdir() + + class LocalClient: + def __init__(self): + self.entries = {} + + def request(self, path, payload=None): + if payload is not None: + name = payload["name"] + entry = {"name": name, "owner": "canopy", + "repository_id": str(uuid.uuid4())} + self.entries[name] = entry + remote = remotes / "canopy" / f"{name}.git" + remote.parent.mkdir(exist_ok=True) + benchmark.git("init", "--bare", "-b", "main", str(remote), + cwd=root, token="fixture-token") + return 200, json.dumps(entry).encode() + name = path.rsplit("/", 1)[-1] + return 200, json.dumps(self.entries[name]).encode() + + client = LocalClient() + manifest = root / "manifest.json" + args = SimpleNamespace(manifest=manifest, work_dir=root / "seed", seed=42, + repositories=2, populated=1, incremental_fixture=True, + base_url=remotes.as_uri()) + self.assertEqual(benchmark.seed(args, client, "fixture-token")["populated"], 1) + entries = benchmark.corpus(manifest)["repositories"] + populated = [entry for entry in entries if entry["commit"] is not None] + self.assertEqual(len(populated), 1) + self.assertNotEqual(populated[0]["base_commit"], populated[0]["commit"]) + check = SimpleNamespace(manifest=manifest, work_dir=root / "verified", + base_url=remotes.as_uri()) + self.assertEqual(benchmark.verify(check, client, "fixture-token") + ["git_v0_v2_samples"], 1) + baseline_manifest = root / "baseline.json" + baseline = SimpleNamespace(manifest=baseline_manifest, work_dir=root / "baseline-seed", + seed=42, repositories=1, populated=1, incremental_fixture=False, + base_url=remotes.as_uri()) + benchmark.seed(baseline, client, "fixture-token") + self.assertEqual(benchmark.corpus(baseline_manifest)["version"], 1) + check = SimpleNamespace(manifest=baseline_manifest, + work_dir=root / "baseline-verified", + base_url=remotes.as_uri()) + self.assertEqual(benchmark.verify(check, client, "fixture-token") + ["git_v0_v2_samples"], 1) + def test_overload_has_no_hidden_retries_or_missing_arrivals(self): server = ThreadingHTTPServer(("127.0.0.1", 0), Handler) server.guard = threading.Lock() @@ -198,6 +248,31 @@ def test_stock_git_clone_and_cold_fetch_validate_populated_repositories(self): token="fixture-token").splitlines() self.assertEqual(refs, [report["push_commit"]] * 2) self.assertNotIn("fixture-token", args.output.read_text()) + entry["base_commit"] = entry["commit"] + benchmark.git("branch", "benchmark-base", cwd=source, token="fixture-token") + increment = b"new commit for an incremental transfer\n" + (source / "incremental.txt").write_bytes(increment) + benchmark.git("add", "incremental.txt", cwd=source, token="fixture-token") + benchmark.git("commit", "-m", "Incremental", cwd=source, token="fixture-token") + entry["commit"] = benchmark.git("rev-parse", "HEAD", cwd=source, + token="fixture-token") + entry["incremental_sha256"] = hashlib.sha256(increment).hexdigest() + benchmark.git("push", str(remote), "HEAD:refs/heads/main", + "benchmark-base:refs/heads/benchmark-base", cwd=source, + token="fixture-token") + manifest.write_text(json.dumps({"version": 2, "complete": True, + "requested_repositories": 1, "repositories": [entry]})) + for operation in ("incremental_fetch", "incremental_pull"): + args = SimpleNamespace(manifest=manifest, active_repositories=1, seed=42, + duration=1, rate=2, concurrency=2, distribution="uniform", + operation=operation, timeout=2, git_timeout=30, + work_dir=root / f"{operation}-scratch", output=root / f"{operation}.json") + report = benchmark.measure(args, ingresses, "fixture-token") + self.assertEqual(report["outcomes"], {"ok": 2}) + self.assertIsNotNone(report["incremental_client_setup_seconds"]) + self.assertEqual([item["outcomes"] for item in report["ingresses"]], + [{"ok": 1}, {"ok": 1}]) + self.assertNotIn("fixture-token", args.output.read_text()) if __name__ == "__main__": From 04403087f0363a860a8ad9ef8d15e440a8b1dea2 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 16:22:51 -0700 Subject: [PATCH 06/12] Benchmark bounded LFS transfers across gateways --- docs/performance-plan.md | 41 ++++++++- scripts/benchmark_repositories.py | 120 +++++++++++++++++++++++-- scripts/test_benchmark_repositories.py | 113 +++++++++++++++++++++++ 3 files changed, 264 insertions(+), 10 deletions(-) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index 58a0287..edfcedc 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -256,8 +256,9 @@ cold repository converge on one serving entry. All pins, denied/lost release, cleanup failure, cancellation and shutdown tests continue to pass. The lock isolation part is implemented. The automated bounded-container run still measures metadata and Git v2 discovery, with stock Git sampling and full identity recovery -checks. The standalone driver can schedule clone, cold/incremental fetch, pull and -unique-ref push, but those workloads have not been qualified at target scale. +checks. The standalone driver can schedule clone, cold/incremental fetch, pull, +unique-ref push and direct-basic LFS transfers, but those workloads have not +been qualified at target scale. Transition queue/service timings are logged. Comprehensive resource metrics and independent load-generator deployment remain open. @@ -509,6 +510,42 @@ transfers, not large histories or cold-cache startup. Client-side shared clone setup is included in each attempt's latency. Report node and generator resource usage alongside the latency distribution. +For LFS basic-transfer load, seed a separate disposable corpus with a declared +body size, then run both read and write workloads through the two ingresses: + +```bash +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --manifest /path/to/canopy-lfs-corpus.json \ + seed --repositories 10000 --populated 100 --lfs-fixture-bytes 4194304 \ + --work-dir /path/to/canopy-lfs-seed +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-lfs-corpus.json \ + run --active-repositories 100 --operation lfs_download \ + --rate 2 --duration 120 --concurrency 8 \ + --output /path/to/canopy-lfs-download.json +python3 -B scripts/benchmark_repositories.py \ + --base-url http://127.0.0.1:8080 \ + --additional-base-url http://127.0.0.1:8081 \ + --manifest /path/to/canopy-lfs-corpus.json \ + run --active-repositories 1000 --operation lfs_upload \ + --lfs-bytes 1048576 --rate 2 --duration 120 --concurrency 8 \ + --output /path/to/canopy-lfs-upload.json +``` + +Seed stores one LFS object per populated repository and `verify` checks its +size and SHA-256. Download attempts stream and hash the full response. Upload +attempts use unique object IDs; their IDs are in the sample file and the +objects remain in the disposable corpus. The driver uses Canopy's direct basic +PUT/GET endpoints, not a stock `git-lfs` batch/checkout flow. It limits fixture +and per-upload bodies to 16 MiB and nominal concurrent payload bytes to +256 MiB; that is not a process-memory ceiling. The separate large-transfer +qualification covers multi-GiB objects. As with +Git, client CPU/disk and network placement must be reported to attribute +throughput. + Run `--operation refs` for Git v2 discovery or `--distribution skewed` for 90% of requests to the selected working set's first tenth. A fixed seed determines working-set selection and offered arrivals. The driver never calls a working set diff --git a/scripts/benchmark_repositories.py b/scripts/benchmark_repositories.py index 58ae959..f46c337 100644 --- a/scripts/benchmark_repositories.py +++ b/scripts/benchmark_repositories.py @@ -2,10 +2,10 @@ """Seed disposable repositories and measure scheduled HTTP and stock-Git work. Use CANOPY_GIT_TOKEN for authentication. Reports contain no credentials. This -measures metadata, Git v2 discovery, clone, fetch, pull or unique-ref push; -LFS throughput needs its own qualification. Seed and verify use stock Git for -a declared corpus sample. Incremental workloads require an opt-in two-commit -corpus. +measures metadata, Git v2 discovery, clone, fetch, pull, unique-ref push or +direct-basic LFS transfers. Seed and verify use stock Git for a declared corpus +sample. Incremental workloads require an opt-in two-commit corpus. Production +throughput still needs separate qualification. """ import argparse @@ -41,13 +41,17 @@ def __init__(self, base, token, timeout): self.connections = [] self.lock = threading.Lock() - def request(self, path, payload=None, *, git=False, request_id=None): + def connection(self): connection = getattr(self.local, "connection", None) if connection is None: connection = self.connection_type(self.host, self.port, timeout=self.timeout) self.local.connection = connection with self.lock: self.connections.append(connection) + return connection + + def request(self, path, payload=None, *, git=False, request_id=None): + connection = self.connection() headers = {"Authorization": f"Bearer {self.token}", "Content-Type": "application/json"} if request_id is not None: headers["X-Request-ID"] = request_id @@ -66,6 +70,47 @@ def request(self, path, payload=None, *, git=False, request_id=None): connection.close() raise + def lfs_put(self, path, content, request_id=None): + connection = self.connection() + headers = {"Authorization": f"Bearer {self.token}", + "Content-Type": "application/octet-stream", + "Content-Length": str(len(content))} + if request_id is not None: + headers["X-Request-ID"] = request_id + try: + connection.request("PUT", self.prefix + path, body=content, headers=headers) + response = connection.getresponse() + if len(response.read(2 * 1024 * 1024 + 1)) > 2 * 1024 * 1024: + raise ValueError("LFS upload response exceeds 2 MiB") + return response.status + except BaseException: + connection.close() + raise + + def lfs_get(self, path, expected_size, request_id=None): + connection = self.connection() + headers = {"Authorization": f"Bearer {self.token}"} + if request_id is not None: + headers["X-Request-ID"] = request_id + try: + connection.request("GET", self.prefix + path, headers=headers) + response = connection.getresponse() + if response.status != 200: + if len(response.read(2 * 1024 * 1024 + 1)) > 2 * 1024 * 1024: + raise ValueError("LFS error response exceeds 2 MiB") + return response.status, 0, None + digest = hashlib.sha256() + size = 0 + while chunk := response.read(1024 * 1024): + size += len(chunk) + if size > expected_size: + raise ValueError("LFS response exceeds declared size") + digest.update(chunk) + return response.status, size, digest.hexdigest() + except BaseException: + connection.close() + raise + def close(self): for connection in self.connections: connection.close() @@ -148,6 +193,15 @@ def seed(args, client, token): git("push", url, *refs, cwd=local, token=token) record["commit"] = git("rev-parse", "HEAD", cwd=local, token=token) record["readme_sha256"] = hashlib.sha256(content).hexdigest() + if args.lfs_fixture_bytes: + lfs_body = hashlib.shake_256( + f"{prefix}/{name}/{args.seed}".encode()).digest(args.lfs_fixture_bytes) + lfs_oid = hashlib.sha256(lfs_body).hexdigest() + path = f"/{entry['owner']}/{name}.git/info/lfs/objects/{lfs_oid}" + if client.lfs_put(path, lfs_body) != 200: + raise RuntimeError("LFS seed upload failed") + record["lfs_oid"] = lfs_oid + record["lfs_size"] = len(lfs_body) if (index + 1) % 25 == 0: save(args.manifest, manifest) print(f"seeded {index + 1}/{args.repositories} repositories", flush=True) @@ -168,6 +222,12 @@ def corpus(path): if any(not re.fullmatch(r"[A-Za-z0-9_-]+", entry[key]) for key in ("name", "owner")): raise ValueError("manifest contains invalid repository names") uuid.UUID(entry["repository_id"]) + if "lfs_oid" in entry: + oid, size = entry["lfs_oid"], entry.get("lfs_size") + if not isinstance(oid, str) or not re.fullmatch(r"[0-9a-f]{64}", oid): + raise ValueError("corpus has an invalid LFS object ID") + if not isinstance(size, int) or isinstance(size, bool) or size < 1: + raise ValueError("corpus has an invalid LFS object size") if manifest["version"] == 2 and entry["commit"] is not None: base_commit = entry.get("base_commit") digest = entry.get("incremental_sha256") @@ -190,6 +250,11 @@ def verify(args, client, token): print(f"verified {index + 1}/{len(manifest['repositories'])} identities", flush=True) if entry["commit"] is None: continue + if "lfs_oid" in entry: + path = f"/{entry['owner']}/{entry['name']}.git/info/lfs/objects/{entry['lfs_oid']}" + status, size, digest = client.lfs_get(path, entry["lfs_size"]) + if (status, size, digest) != (200, entry["lfs_size"], entry["lfs_oid"]): + raise RuntimeError("restored LFS bytes differ from manifest") for protocol in ("0", "2"): clone = args.work_dir / f"{entry['name']}-v{protocol}" url = f"{args.base_url.rstrip('/')}/{entry['owner']}/{entry['name']}.git" @@ -293,7 +358,14 @@ def measure(args, client, token): git_read = args.operation in ("clone", "cold_fetch") or incremental git_write = args.operation == "push_branch" git_operation = git_read or git_write - eligible = ([entry for entry in manifest["repositories"] if entry.get("base_commit") is not None] + lfs_download = args.operation == "lfs_download" + lfs_upload = args.operation == "lfs_upload" + if lfs_upload and (args.lfs_bytes == 0 or + args.lfs_bytes * args.concurrency > 256 * 1024 * 1024): + raise ValueError("LFS uploads require positive bytes and at most 256 MiB in-flight payloads") + eligible = ([entry for entry in manifest["repositories"] if entry.get("lfs_oid") is not None] + if lfs_download else + [entry for entry in manifest["repositories"] if entry.get("base_commit") is not None] if incremental else [entry for entry in manifest["repositories"] if entry["commit"] is not None] if git_read else manifest["repositories"]) @@ -310,6 +382,9 @@ def measure(args, client, token): raise ValueError("stock-Git runs require --work-dir") args.work_dir.mkdir(parents=True, exist_ok=False) push_run_id = uuid.uuid4().hex if git_write else None + lfs_run_id = uuid.uuid4().hex if lfs_upload else None + lfs_tail = (hashlib.shake_256(lfs_run_id.encode()).digest(max(0, args.lfs_bytes - 32)) + if lfs_upload else None) source = args.work_dir / "source" if git_write else None push_commit = None if git_write: @@ -357,6 +432,7 @@ def execute(sequence, entry, scheduled): ingress = clients[ingress_index] dispatched = time.monotonic() result = "transport_error" + uploaded_oid = None try: if args.operation == "metadata": status, body = ingress.request(f"/api/repositories/{entry['name']}", request_id=request_id) @@ -370,11 +446,24 @@ def execute(sequence, entry, scheduled): valid = git_transfer(args.operation, url, entry, token, args.work_dir, args.git_timeout, request_id, templates.get(entry["repository_id"])) - else: + elif git_write: url = f"{ingress.base_url}/{entry['owner']}/{entry['name']}.git" reference = f"refs/heads/canopy-benchmark/{push_run_id}/{sequence:07d}" push_branch(url, source, reference, token, args.git_timeout, request_id) valid = True + elif lfs_download: + path = f"/{entry['owner']}/{entry['name']}.git/info/lfs/objects/{entry['lfs_oid']}" + status, size, digest = ingress.lfs_get(path, entry["lfs_size"], request_id) + valid = (status, size, digest) == (200, entry["lfs_size"], entry["lfs_oid"]) + elif lfs_upload: + marker = hashlib.sha256(f"{lfs_run_id}:{sequence}".encode()).digest() + body = (marker + lfs_tail)[:args.lfs_bytes] + uploaded_oid = hashlib.sha256(body).hexdigest() + path = f"/{entry['owner']}/{entry['name']}.git/info/lfs/objects/{uploaded_oid}" + status = ingress.lfs_put(path, body, request_id) + valid = status == 200 + else: + raise ValueError("unsupported benchmark operation") result = ("ok" if valid else (f"http_{status}" if not git_operation and status != 200 else "invalid_response")) except subprocess.TimeoutExpired: @@ -386,6 +475,7 @@ def execute(sequence, entry, scheduled): finally: finished = time.monotonic() record({"sequence": sequence, "request_id": request_id, "repository_id": entry["repository_id"], "ingress_index": ingress_index, "result": result, + "lfs_oid": uploaded_oid, "elapsed_ms": (finished - scheduled) * 1000, "service_ms": (finished - dispatched) * 1000, "dispatch_delay_ms": (dispatched - scheduled) * 1000}) @@ -417,6 +507,8 @@ def execute(sequence, entry, scheduled): "concurrency": args.concurrency, "request_timeout_seconds": args.timeout, "git_timeout_seconds": args.git_timeout if git_operation else None, "push_run_id": push_run_id, "push_commit": push_commit, + "lfs_run_id": lfs_run_id, + "lfs_size_bytes": args.lfs_bytes if lfs_upload else None, "scheduled": total, "outcomes": dict(counts), "ingresses": [{"index": index, "outcomes": dict(outcomes), "scheduled_latency_ms": percentiles(ingress_latencies[index]), @@ -440,6 +532,13 @@ def positive(value): return parsed +def bounded_lfs_bytes(value): + parsed = int(value) + if not 0 <= parsed <= 16 * 1024 * 1024: + raise argparse.ArgumentTypeError("LFS benchmark body must be 0 to 16 MiB") + return parsed + + def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--base-url", required=True) @@ -454,6 +553,8 @@ def main(): create.add_argument("--populated", type=positive, default=3) create.add_argument("--incremental-fixture", action="store_true", help="seed a second commit and benchmark-base ref for incremental fetch/pull") + create.add_argument("--lfs-fixture-bytes", type=bounded_lfs_bytes, default=0, + help="also upload one direct-basic LFS object per populated repository") create.add_argument("--work-dir", type=Path, required=True) check = commands.add_parser("verify") check.add_argument("--work-dir", type=Path, required=True) @@ -461,13 +562,16 @@ def main(): run.add_argument("--active-repositories", type=positive, required=True) run.add_argument("--distribution", choices=("uniform", "skewed"), default="uniform") run.add_argument("--operation", choices=("metadata", "refs", "clone", "cold_fetch", - "incremental_fetch", "incremental_pull", "push_branch"), + "incremental_fetch", "incremental_pull", "push_branch", + "lfs_download", "lfs_upload"), default="metadata") run.add_argument("--rate", type=positive, default=20) run.add_argument("--duration", type=positive, default=30) run.add_argument("--concurrency", type=positive, default=32) run.add_argument("--work-dir", type=Path, help="new client directory for stock-Git operations") run.add_argument("--git-timeout", type=positive, default=120) + run.add_argument("--lfs-bytes", type=bounded_lfs_bytes, default=1024 * 1024, + help="unique LFS upload bytes per scheduled arrival, at most 16 MiB") run.add_argument("--output", type=Path, required=True) args = parser.parse_args() token = os.environ.get("CANOPY_GIT_TOKEN") diff --git a/scripts/test_benchmark_repositories.py b/scripts/test_benchmark_repositories.py index f80af18..2ab98a8 100644 --- a/scripts/test_benchmark_repositories.py +++ b/scripts/test_benchmark_repositories.py @@ -39,6 +39,34 @@ def do_GET(self): self.server.active -= 1 +class LfsHandler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" + + def log_message(self, *_args): + pass + + def do_PUT(self): + size = int(self.headers["Content-Length"]) + body = self.rfile.read(size) + oid = self.path.rsplit("/", 1)[-1] + status = 200 if hashlib.sha256(body).hexdigest() == oid else 422 + if status == 200: + with self.server.guard: + self.server.blobs[self.path] = body + self.send_response(status) + self.send_header("Content-Length", "0") + self.end_headers() + + def do_GET(self): + with self.server.guard: + body = self.server.blobs.get(self.path) + self.send_response(200 if body is not None else 404) + self.send_header("Content-Length", str(len(body) if body is not None else 0)) + self.end_headers() + if body is not None: + self.wfile.write(body) + + class ScheduledLoad(unittest.TestCase): def test_incremental_seed_and_verify_preserve_both_commits(self): with tempfile.TemporaryDirectory() as directory: @@ -49,6 +77,7 @@ def test_incremental_seed_and_verify_preserve_both_commits(self): class LocalClient: def __init__(self): self.entries = {} + self.lfs = {} def request(self, path, payload=None): if payload is not None: @@ -64,16 +93,30 @@ def request(self, path, payload=None): name = path.rsplit("/", 1)[-1] return 200, json.dumps(self.entries[name]).encode() + def lfs_put(self, path, content, request_id=None): + if hashlib.sha256(content).hexdigest() != path.rsplit("/", 1)[-1]: + return 422 + self.lfs[path] = content + return 200 + + def lfs_get(self, path, expected_size, request_id=None): + body = self.lfs.get(path) + if body is None: + return 404, 0, None + return 200, len(body), hashlib.sha256(body).hexdigest() + client = LocalClient() manifest = root / "manifest.json" args = SimpleNamespace(manifest=manifest, work_dir=root / "seed", seed=42, repositories=2, populated=1, incremental_fixture=True, + lfs_fixture_bytes=128, base_url=remotes.as_uri()) self.assertEqual(benchmark.seed(args, client, "fixture-token")["populated"], 1) entries = benchmark.corpus(manifest)["repositories"] populated = [entry for entry in entries if entry["commit"] is not None] self.assertEqual(len(populated), 1) self.assertNotEqual(populated[0]["base_commit"], populated[0]["commit"]) + self.assertEqual(populated[0]["lfs_size"], 128) check = SimpleNamespace(manifest=manifest, work_dir=root / "verified", base_url=remotes.as_uri()) self.assertEqual(benchmark.verify(check, client, "fixture-token") @@ -81,6 +124,7 @@ def request(self, path, payload=None): baseline_manifest = root / "baseline.json" baseline = SimpleNamespace(manifest=baseline_manifest, work_dir=root / "baseline-seed", seed=42, repositories=1, populated=1, incremental_fixture=False, + lfs_fixture_bytes=0, base_url=remotes.as_uri()) benchmark.seed(baseline, client, "fixture-token") self.assertEqual(benchmark.corpus(baseline_manifest)["version"], 1) @@ -198,6 +242,75 @@ def test_stock_git_sends_authorization_and_request_id_headers(self): server.server_close() thread.join() + def test_lfs_upload_and_streamed_download_across_ingresses(self): + with tempfile.TemporaryDirectory() as directory: + root = Path(directory) + body = b"LFS fixture\n" * (3 * 1024 * 1024 // len(b"LFS fixture\n")) + oid = hashlib.sha256(body).hexdigest() + path = f"/canopy/fixture.git/info/lfs/objects/{oid}" + blobs = {path: body} + guard = threading.Lock() + servers = [ThreadingHTTPServer(("127.0.0.1", 0), LfsHandler) for _ in range(2)] + threads = [] + clients = [] + for server in servers: + server.blobs = blobs + server.guard = guard + thread = threading.Thread(target=server.serve_forever, daemon=True) + thread.start() + threads.append(thread) + clients.append(benchmark.Client(f"http://127.0.0.1:{server.server_port}", + "fixture-token", 5)) + try: + manifest = root / "manifest.json" + manifest.write_text(json.dumps({"version": 1, "complete": True, + "requested_repositories": 1, "repositories": [ + {"name": "fixture", "owner": "canopy", + "repository_id": str(uuid.uuid4()), "commit": None, + "lfs_oid": oid, "lfs_size": len(body)}]})) + args = SimpleNamespace(manifest=manifest, active_repositories=1, seed=42, + duration=1, rate=2, concurrency=2, distribution="uniform", + operation="lfs_download", timeout=5, output=root / "download.json") + report = benchmark.measure(args, clients, "fixture-token") + self.assertEqual(report["outcomes"], {"ok": 2}) + self.assertEqual([item["outcomes"] for item in report["ingresses"]], + [{"ok": 1}, {"ok": 1}]) + with guard: + blobs[path] = b"X" + body[1:] + args.output = root / "corrupt-download.json" + report = benchmark.measure(args, clients, "fixture-token") + self.assertEqual(report["outcomes"], {"invalid_response": 2}) + self.assertEqual(report["failed_arrivals"], 2) + with guard: + blobs[path] = body + args.operation = "lfs_upload" + args.lfs_bytes = 1024 + args.output = root / "upload.json" + report = benchmark.measure(args, clients, "fixture-token") + self.assertEqual(report["outcomes"], {"ok": 2}) + samples = [json.loads(line) for line in + (root / "upload.samples.jsonl").read_text().splitlines()] + self.assertEqual(len({sample["lfs_oid"] for sample in samples}), 2) + for sample in samples: + uploaded = blobs[f"/canopy/fixture.git/info/lfs/objects/{sample['lfs_oid']}"] + self.assertEqual(len(uploaded), 1024) + self.assertEqual(hashlib.sha256(uploaded).hexdigest(), sample["lfs_oid"]) + self.assertNotIn("fixture-token", args.output.read_text()) + args.lfs_bytes = 16 * 1024 * 1024 + args.concurrency = 32 + args.output = root / "too-large.json" + with self.assertRaises(ValueError): + benchmark.measure(args, clients, "fixture-token") + self.assertFalse(args.output.exists()) + finally: + for client in clients: + client.close() + for server in servers: + server.shutdown() + server.server_close() + for thread in threads: + thread.join() + def test_stock_git_clone_and_cold_fetch_validate_populated_repositories(self): with tempfile.TemporaryDirectory() as directory: root = Path(directory) From 2ab4b2c7f2c139006a5db6dd3150034c97b02cfb Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 16:53:06 -0700 Subject: [PATCH 07/12] Pin Cellule renewal scheduler for dense repository ownership --- Cargo.lock | 12 ++++++------ Cargo.toml | 10 +++++----- docs/performance-plan.md | 25 +++++++++++++++++++------ 3 files changed, 30 insertions(+), 17 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 4c50efc..f4b8e49 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -419,7 +419,7 @@ dependencies = [ [[package]] name = "cellule-app" version = "0.1.0" -source = "git+https://github.com/crabbuild/cellule.git?rev=a28de7bc09ce36d87e642adc4f4b6be50d6fcb69#a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" +source = "git+https://github.com/crabbuild/cellule.git?rev=47a302b79962ee16c698e121315cf4e85ec49549#47a302b79962ee16c698e121315cf4e85ec49549" dependencies = [ "blake3", "cellule-runtime", @@ -428,7 +428,7 @@ dependencies = [ [[package]] name = "cellule-host" version = "0.1.0" -source = "git+https://github.com/crabbuild/cellule.git?rev=a28de7bc09ce36d87e642adc4f4b6be50d6fcb69#a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" +source = "git+https://github.com/crabbuild/cellule.git?rev=47a302b79962ee16c698e121315cf4e85ec49549#47a302b79962ee16c698e121315cf4e85ec49549" dependencies = [ "cellule-app", "cellule-runtime", @@ -442,7 +442,7 @@ dependencies = [ [[package]] name = "cellule-ltx" version = "0.1.0" -source = "git+https://github.com/crabbuild/cellule.git?rev=a28de7bc09ce36d87e642adc4f4b6be50d6fcb69#a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" +source = "git+https://github.com/crabbuild/cellule.git?rev=47a302b79962ee16c698e121315cf4e85ec49549#47a302b79962ee16c698e121315cf4e85ec49549" dependencies = [ "async-trait", "blake3", @@ -464,7 +464,7 @@ dependencies = [ [[package]] name = "cellule-runtime" version = "0.1.0" -source = "git+https://github.com/crabbuild/cellule.git?rev=a28de7bc09ce36d87e642adc4f4b6be50d6fcb69#a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" +source = "git+https://github.com/crabbuild/cellule.git?rev=47a302b79962ee16c698e121315cf4e85ec49549#47a302b79962ee16c698e121315cf4e85ec49549" dependencies = [ "blake3", "bytes", @@ -490,7 +490,7 @@ dependencies = [ [[package]] name = "cellule-store" version = "0.1.0" -source = "git+https://github.com/crabbuild/cellule.git?rev=a28de7bc09ce36d87e642adc4f4b6be50d6fcb69#a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" +source = "git+https://github.com/crabbuild/cellule.git?rev=47a302b79962ee16c698e121315cf4e85ec49549#47a302b79962ee16c698e121315cf4e85ec49549" dependencies = [ "async-trait", "blake3", @@ -513,7 +513,7 @@ dependencies = [ [[package]] name = "cellule-types" version = "0.1.0" -source = "git+https://github.com/crabbuild/cellule.git?rev=a28de7bc09ce36d87e642adc4f4b6be50d6fcb69#a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" +source = "git+https://github.com/crabbuild/cellule.git?rev=47a302b79962ee16c698e121315cf4e85ec49549#47a302b79962ee16c698e121315cf4e85ec49549" dependencies = [ "schemars", "serde", diff --git a/Cargo.toml b/Cargo.toml index 558bc49..49b2a88 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,11 +13,11 @@ axum = "0.8.9" base64 = "0.22" blake3 = "1.8" bytes = "1.11" -cellule-app = { git = "https://github.com/crabbuild/cellule.git", rev = "a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" } -cellule-host = { git = "https://github.com/crabbuild/cellule.git", rev = "a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" } -cellule-ltx = { git = "https://github.com/crabbuild/cellule.git", rev = "a28de7bc09ce36d87e642adc4f4b6be50d6fcb69", features = ["replica"] } -cellule-runtime = { git = "https://github.com/crabbuild/cellule.git", rev = "a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" } -cellule-store = { git = "https://github.com/crabbuild/cellule.git", rev = "a28de7bc09ce36d87e642adc4f4b6be50d6fcb69" } +cellule-app = { git = "https://github.com/crabbuild/cellule.git", rev = "47a302b79962ee16c698e121315cf4e85ec49549" } +cellule-host = { git = "https://github.com/crabbuild/cellule.git", rev = "47a302b79962ee16c698e121315cf4e85ec49549" } +cellule-ltx = { git = "https://github.com/crabbuild/cellule.git", rev = "47a302b79962ee16c698e121315cf4e85ec49549", features = ["replica"] } +cellule-runtime = { git = "https://github.com/crabbuild/cellule.git", rev = "47a302b79962ee16c698e121315cf4e85ec49549" } +cellule-store = { git = "https://github.com/crabbuild/cellule.git", rev = "47a302b79962ee16c698e121315cf4e85ec49549" } ed25519-dalek = "2" flate2 = "1.1" futures-core = "0.3" diff --git a/docs/performance-plan.md b/docs/performance-plan.md index edfcedc..a9463ba 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -1,6 +1,6 @@ # Measure repository density and latency -Use this plan to design capacity work and interpret Canopy benchmark results. It separates demonstrated behavior from proposed targets. The current Cellule dependency is `a28de7bc09ce36d87e642adc4f4b6be50d6fcb69`; earlier runs below retain their original pins and do not establish this build's density or latency. +Use this plan to design capacity work and interpret Canopy benchmark results. It separates demonstrated behavior from proposed targets. The current Cellule dependency is `47a302b79962ee16c698e121315cf4e85ec49549`; earlier runs below retain their original pins and do not establish this build's density or latency. ## Read the result before the target @@ -14,6 +14,17 @@ Use this plan to design capacity work and interpret Canopy benchmark results. It The [implementation order](#implementation-order-and-acceptance) defines work still needed. The [measurement history](#measurement-history) records the revision, hardware, provider and workload for individual runs. Compare those four inputs before combining numbers from different sections. +The current Cellule pin dispatches due owner renewals oldest-first and refills +the bounded 32-task renewal window when I/O completes. A 100-ms scan rebuilds +pending candidates, removing stale generations and departed Cells. This removes +the earlier 320-starts/s *scheduler* ceiling; it does not lower the one-control- +update-per-active-Cell cost or prove that a provider can sustain the required +update rate. Repeat the real-store idle coverage gate before raising the active +Cell limit or claiming 1,000- or 10,000-Cell residency. This dependency and +lockfile change also changes Canopy's compiled release digest; test against a +fresh store prefix or use the documented maintenance upgrade path for an +existing deployment. + ## Recorded large-transfer qualification The release-mode `scripts/qualify_size.py --docker-volume --release` gate passed @@ -1654,12 +1665,14 @@ The dispatch ceiling is source-backed; exact starvation causality still needs a controlled scheduler regression. The existing `idle_owner_progress_is_renewed_without_a_per_cell_task` test covers one Cell. -Before raising active density, implement and qualify these runtime changes: +The current Cellule pin addresses the dispatch-order and refill mechanism with +a sorted scan every 100 ms. Before raising active density, qualify the change: -1. Dispatch the oldest eligible overdue renewal first, with bounded concurrency - and capacity refill after completion. Avoid a full Cell scan per completion; - use a deadline index with generation checks and bounded stale-entry cleanup. -2. Preserve coordination admission, publisher exclusivity, node lease checks, +1. Measure scan and sort CPU at 1,000 through 10,000 active Cells. The current + pending vector is rebuilt each tick and bounded by active count; use a + deadline index only if measured scan cost warrants the extra invalidation + and generation bookkeeping. +2. Verify coordination admission, publisher exclusivity, node lease checks, effect/generation fencing and shutdown behavior. Foreground publication, compaction, transfer and release share these invariants. 3. Prove progress above the old dispatch ceiling, including slow storage, From 1ee42edbffa1880ba73759f2c991ff19b5488471 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 17:17:47 -0700 Subject: [PATCH 08/12] Record renewal density evidence and add memory diagnostic mode --- docs/performance-plan.md | 37 ++++++++++++++++++++++++++++++++++++- examples/benchmark_idle.rs | 12 ++++++++++-- 2 files changed, 46 insertions(+), 3 deletions(-) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index a9463ba..fdde2e3 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -10,7 +10,7 @@ Use this plan to design capacity work and interpret Canopy benchmark results. It | Can a bounded Linux node recover 1,000 repository identities? | Yes, in a mostly empty SQL-only corpus; 997 repositories were empty | [Bounded Linux density qualification](#bounded-linux-density-qualification) | | Did that run meet every warm metadata latency target? | No; some p95 and p99 targets were missed | [Bounded Linux density qualification](#bounded-linux-density-qualification) | | Is 10,000 repositories per node a measured capacity? | No; it is a proposed reference-node target | [Required outcome](#required-outcome) and [qualification rules](#performance-qualification-rules) | -| Is idle ownership proven at 1,000 active Cells? | No; a run on an older Cellule revision missed renewal coverage for nine Cells, and the current pin needs a repeat | [Observed idle renewal ceiling](#observed-idle-renewal-ceiling) | +| Is idle ownership proven at 1,000 active Cells? | Not on a real store; the current pin passed 1,000 in-memory Cells, but its RustFS run fenced owners beyond 800 | [Current scheduler and provider observations](#current-pin-scheduler-and-provider-observations) | The [implementation order](#implementation-order-and-acceptance) defines work still needed. The [measurement history](#measurement-history) records the revision, hardware, provider and workload for individual runs. Compare those four inputs before combining numbers from different sections. @@ -1608,6 +1608,15 @@ cargo run --release --locked --example benchmark_idle -- \ 0,100,500,1000 30 ``` +For a scheduler-only diagnostic, use `memory:///` with a different new output +directory. The report labels this provider as in-memory; a passing result +cannot replace the S3-compatible gate or establish provider throughput. + +```bash +cargo run --locked --example benchmark_idle -- \ + memory:/// /tmp/canopy-idle-memory-diagnostic 100,500,1000 30 +``` + Each phase waits four seconds after foreground work, then measures without HTTP requests for the requested interval (10–60 seconds). The report records successful conditional Cell-control updates, unchanged-root updates, root changes, updated @@ -1665,6 +1674,32 @@ The dispatch ceiling is source-backed; exact starvation causality still needs a controlled scheduler regression. The existing `idle_owner_progress_is_renewed_without_a_per_cell_task` test covers one Cell. +#### Current pin: scheduler and provider observations + +The current Cellule pin is `47a302b`. A debug Canopy example on shared macOS +arm64 measured these SQL-only, otherwise idle windows; each active count also +includes the Directory Cell. The in-memory run passed its final graceful drain, +fresh-workspace restart, released-only window, and three identity restorations. +The beta.8 RustFS container was capped at 2 CPUs and 4 GiB. + +| Store | Active repository Cells | Window | Control updates/s | Distinct Cells updated | Result | +| --- | ---: | ---: | ---: | ---: | --- | +| In-memory | 100 | 30 s | 32.629 | 101/101 | Passed | +| In-memory | 500 | 30 s | 161.616 | 501/501 | Passed | +| In-memory | 1,000 | 30 s | 326.830 | 1,001/1,001 | Passed | +| RustFS beta.8 | 100 | 10 s | 32.793 | 101/101 | Passed | +| RustFS beta.8 | 500 | 10 s | 163.759 | 501/501 | Passed | + +All listed windows had zero failed PUTs and no root changes. The in-memory +1,000-Cell window exceeds the former 320-starts/s tick ceiling, isolating the +scheduler improvement. It does **not** qualify real-store density. During the +beta.8 run, repository creation stopped advancing at 812 after many Cell +renewals returned `Fenced`; the process did not complete the 1,000-Cell window +or its graceful shutdown and was terminated for diagnosis. A basic provider +bucket check still responded. This is a failed/incomplete real-store gate, not +proof of a single root cause; node-lease refresh, provider latency and shutdown +drain need investigation before increasing active density. + The current Cellule pin addresses the dispatch-order and refill mechanism with a sorted scan every 100 ms. Before raising active density, qualify the change: diff --git a/examples/benchmark_idle.rs b/examples/benchmark_idle.rs index cf186b0..0c6e850 100644 --- a/examples/benchmark_idle.rs +++ b/examples/benchmark_idle.rs @@ -21,7 +21,9 @@ type Result = std::result::Result; #[derive(Debug, thiserror::Error)] enum Error { - #[error("usage: benchmark_idle =10")] + #[error( + "usage: benchmark_idle =10" + )] Usage, #[error("benchmark invariant failed: {0}")] Invalid(&'static str), @@ -98,7 +100,7 @@ async fn run() -> Result { return Err(Error::Usage); } let url = url::Url::parse(&args[0])?; - if url.scheme() != "s3" + if !matches!(url.scheme(), "s3" | "memory") || !url.username().is_empty() || url.password().is_some() || url.query().is_some() @@ -146,7 +148,13 @@ async fn run() -> Result { token: format!("cnp_{}", hex::encode(Sha256::digest(key))), active_limit, }; + let provider_kind = if url.scheme() == "memory" { + "in-memory scheduler diagnostic; not a real-store qualification" + } else { + "S3-compatible provider" + }; let mut report = json!({"passed":false,"cell_capabilities":"SQL only","requested_counts":counts,"window_seconds":seconds, + "provider_kind":provider_kind, "binary_sha256":hex::encode(Sha256::digest(std::fs::read(std::env::current_exe()?)?)),"store_prefix":fixture.prefix.to_string(), "counting_boundary":"object_store API, not provider HTTP attempts; list/delete/multipart parts are not counted", "windows":[],"seeded_repositories":[],"shutdown_passed":false}); From e57d149370a606b18aeabf923fcf1dd21fefd424 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 17:27:12 -0700 Subject: [PATCH 09/12] Skip empty ref scans in initial mirror publication --- src/refs.rs | 29 +++++++++++++++++++++++++++-- tests/repository_cell.rs | 27 +++++++++++++++++++++++++++ 2 files changed, 54 insertions(+), 2 deletions(-) diff --git a/src/refs.rs b/src/refs.rs index 72e1bbb..0aa4cc1 100644 --- a/src/refs.rs +++ b/src/refs.rs @@ -349,6 +349,24 @@ pub(crate) fn apply_refs( if updates.len() != plan.updates.len() { return Ok(false); } + // A first mirror push can contain thousands of refs. One transactional + // emptiness check avoids a point lookup and namespace scan for every new + // name while preserving the ordinary CAS path for tombstones and live refs. + let emptiness = context.sql(&SqlBatch { + statements: vec![SqlStatement { + sql: "SELECT NOT EXISTS (SELECT 1 FROM refs)".into(), + parameters: Vec::new(), + }], + })?; + let refs_empty = match emptiness + .first() + .and_then(|set| set.rows.first()) + .map(Vec::as_slice) + { + Some([SqlValue::Integer(0)]) => false, + Some([SqlValue::Integer(1)]) => true, + _ => return Err(Error::Command("invalid ref emptiness result")), + }; for update in &plan.updates { if server_owned_ref(&update.name) || !valid_ref_name(&update.name) @@ -362,7 +380,9 @@ pub(crate) fn apply_refs( if update.new_oid.is_none() && update.expected.as_ref().and_then(|old| old.oid).is_none() { return Ok(false); } - if current_ref(context, &update.name)? != update.expected { + if (refs_empty && update.expected.is_some()) + || (!refs_empty && current_ref(context, &update.name)? != update.expected) + { return Ok(false); } } @@ -371,7 +391,7 @@ pub(crate) fn apply_refs( .iter() .filter(|update| update.new_oid.is_some()) { - if existing_namespace_conflict(context, &updates, &update.name)? { + if existing_namespace_conflict(context, &updates, &update.name, refs_empty)? { return Ok(false); } } @@ -474,6 +494,7 @@ fn existing_namespace_conflict( context: &CommandContext<'_, '_>, updates: &BTreeMap<&str, &RefUpdate>, name: &str, + refs_empty: bool, ) -> cellule_runtime::Result { // Planned deletions remove namespace conflicts in this same transaction. // Exact ancestor lookups and indexed descendant pages avoid a full ref scan @@ -482,6 +503,7 @@ fn existing_namespace_conflict( let ancestor = &name[..index]; let live = match updates.get(ancestor) { Some(update) => update.new_oid.is_some(), + None if refs_empty => false, None => current_ref(context, ancestor)?.is_some_and(|state| state.oid.is_some()), }; if live { @@ -496,6 +518,9 @@ fn existing_namespace_conflict( { return Ok(true); } + if refs_empty { + return Ok(false); + } let mut after = prefix.clone(); loop { let result = context.sql(&SqlBatch { diff --git a/tests/repository_cell.rs b/tests/repository_cell.rs index 14bdb43..77f2470 100644 --- a/tests/repository_cell.rs +++ b/tests/repository_cell.rs @@ -225,6 +225,33 @@ async fn repository_cell_publishes_objects_and_refs_atomically() second_commit.as_bytes(), ) .await?; + // The empty-ref fast path still rejects an ancestor and descendant in + // the same atomic plan, with no generation change or partial writes. + let before = repository.refs_page("", None).await?.output; + assert!(before.refs.is_empty()); + assert!(matches!( + repository + .finalize_push( + identity(28), + PushPlan { + actor: "canopy".into(), + updates: ["refs/tags/conflict", "refs/tags/conflict/child"] + .into_iter() + .map(|name| RefUpdate { + name: name.into(), + expected: None, + new_oid: Some(committed.output), + }) + .collect(), + }, + ) + .await, + Err(InvocationError::Rejected(_)) + )); + assert_eq!( + repository.default_branch(None).await?.output.generation, + before.generation + ); let published = repository .finalize_push( MutationIdentity { From 2b51b141e47de21c7228e6122e08cb231d6a40a9 Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 17:48:20 -0700 Subject: [PATCH 10/12] Make cold-residency fault fixtures establish eviction --- tests/multi_server/residency/faults.rs | 47 ++++++++++++++++++++++---- 1 file changed, 40 insertions(+), 7 deletions(-) diff --git a/tests/multi_server/residency/faults.rs b/tests/multi_server/residency/faults.rs index c5f3528..18218cb 100644 --- a/tests/multi_server/residency/faults.rs +++ b/tests/multi_server/residency/faults.rs @@ -269,6 +269,40 @@ impl Fixture { run_git(Some(&clone), &["fsck", "--full"]).await?; Ok(()) } + + async fn make_original_cold(&self) -> Result { + // A freshly pushed Cell can retain publication work briefly, so the + // first capacity eviction need not choose it. Keep the comparison + // repository warm and request bounded additional identities until the + // fixture actually reaches the cold-restore precondition. + for attempt in 0..20 { + if attempt > 0 { + tokio::time::sleep(Duration::from_millis(250)).await; + self.client + .get(format!("http://{}/api/repositories/third", self.address)) + .bearer_auth("local-test-token") + .send() + .await? + .error_for_status()?; + } + let name = if attempt == 0 { + "fourth".to_owned() + } else { + format!("eviction-{attempt}") + }; + create(&self.client, self.address, &name).await?; + if !self.repository_dir.exists() { + self.client + .get(format!("http://{}/api/repositories/third", self.address)) + .bearer_auth("local-test-token") + .send() + .await? + .error_for_status()?; + return Ok(()); + } + } + Err("original repository never became cold".into()) + } } #[tokio::test(flavor = "multi_thread")] @@ -416,8 +450,7 @@ async fn paused_release_keeps_other_warm_repositories_available() -> Result { #[tokio::test(flavor = "multi_thread")] async fn paused_cold_activation_keeps_other_warm_repositories_available() -> Result { let fixture = Fixture::new().await?; - create(&fixture.client, fixture.address, "fourth").await?; - assert!(!fixture.repository_dir.exists()); + fixture.make_original_cold().await?; *fixture.store.paused_read.lock().unwrap() = Some( fixture .layout @@ -455,7 +488,9 @@ async fn paused_cold_activation_keeps_other_warm_repositories_available() -> Res impl Fixture { async fn read_warm_repository(&self) -> Result { - timeout(Duration::from_secs(2), async { + // This fault test asserts independence from a paused release/restore, + // not a two-second latency SLO on an oversubscribed test host. + timeout(Duration::from_secs(5), async { let response = self .client .get(format!("http://{}/api/repositories/third", self.address)) @@ -490,8 +525,7 @@ impl Fixture { #[tokio::test(flavor = "multi_thread")] async fn paused_cold_repository_does_not_serialize_other_cold_activations() -> Result { let fixture = Fixture::new().await?; - create(&fixture.client, fixture.address, "fourth").await?; - assert!(!fixture.repository_dir.exists()); + fixture.make_original_cold().await?; *fixture.store.paused_read.lock().unwrap() = Some( fixture .layout @@ -524,8 +558,7 @@ async fn paused_cold_repository_does_not_serialize_other_cold_activations() -> R #[tokio::test(flavor = "multi_thread")] async fn cancelled_cold_activation_retains_its_reserved_slot() -> Result { let fixture = Fixture::new().await?; - create(&fixture.client, fixture.address, "fourth").await?; - assert!(!fixture.repository_dir.exists()); + fixture.make_original_cold().await?; let mut uploads = Vec::new(); let oid = hex::encode(Sha256::digest(b"x")); for name in ["third", "fourth"] { From 97e9ddd25d68d5ed852c52bcbfcdb26c4d6520da Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 17:58:59 -0700 Subject: [PATCH 11/12] Fail large-transfer qualification early on undersized volumes --- docs/performance-plan.md | 6 +++++ scripts/qualify_size.py | 22 ++++++++++++++++++ scripts/test_qualify_size.py | 45 ++++++++++++++++++++++++++++++++++++ 3 files changed, 73 insertions(+) create mode 100644 scripts/test_qualify_size.py diff --git a/docs/performance-plan.md b/docs/performance-plan.md index fdde2e3..87d40b8 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -41,6 +41,12 @@ preparation hydrated 901 reachable blobs (5,966,921,728 bytes) in about 109 seconds. This is functional size and recovery proof on the local provider; it does not establish production-provider latency or repository density. +The current PR's GitHub-hosted size gate has not repeated that pass. RustFS +returned a write-quorum HTTP 500 while staging the large Git blob, and the +runner then reported `No space left on device`. The gate now checks for at least +40 GiB free on both its scratch and provider volumes before uploading; it still +requires the full non-sparse Git and LFS round trip on a suitable runner. + ## Required outcome A node must serve thousands of repository identities, each with its own diff --git a/scripts/qualify_size.py b/scripts/qualify_size.py index ca87b5b..8f5b811 100644 --- a/scripts/qualify_size.py +++ b/scripts/qualify_size.py @@ -2,17 +2,37 @@ import argparse import os from pathlib import Path +import shutil import subprocess import tempfile import time import uuid +MIN_FREE_BYTES = 40 * 1024**3 + def run(*args, **kwargs): result = subprocess.run(args, check=True, capture_output=True, text=True, **kwargs) return (result.stdout + (result.stderr if args[:2] == ("docker", "logs") else "")).strip() +def require_size_capacity(directory: Path, container: str): + """Check both scratch and provider volumes before the non-sparse 5 GiB gate.""" + scratch_free = shutil.disk_usage(directory).free + output = run("docker", "exec", container, "df", "-B1", "/data") + try: + provider_free = int(output.splitlines()[-1].split()[3]) + except (IndexError, ValueError) as error: + raise RuntimeError("could not read free space on RustFS /data") from error + if min(scratch_free, provider_free) < MIN_FREE_BYTES: + gib = 1024**3 + raise RuntimeError( + "large Git/LFS qualification requires at least 40 GiB free on both " + f"scratch and RustFS /data (scratch={scratch_free / gib:.1f} GiB, " + f"provider={provider_free / gib:.1f} GiB); use a larger runner/volume" + ) + + def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--provider-only", action="store_true") @@ -77,6 +97,8 @@ def main(): time.sleep(1) else: raise RuntimeError("RustFS fixture did not become ready") + if not args.provider_only: + require_size_capacity(directory, name) env["CANOPY_TEST_S3_ENDPOINT"] = endpoint env["CANOPY_TEST_S3_BUCKET"] = "canopy-size" tests = [] if args.provider_only else [ diff --git a/scripts/test_qualify_size.py b/scripts/test_qualify_size.py new file mode 100644 index 0000000..42e534a --- /dev/null +++ b/scripts/test_qualify_size.py @@ -0,0 +1,45 @@ +"""Preflight checks for the large, non-sparse RustFS qualification.""" + +import importlib.util +from pathlib import Path +import unittest +from unittest.mock import patch + + +SPEC = importlib.util.spec_from_file_location( + "qualify_size", Path(__file__).with_name("qualify_size.py") +) +qualify_size = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(qualify_size) + + +class SizeCapacityTests(unittest.TestCase): + def check(self, scratch_gib, provider_gib): + with patch.object(qualify_size.shutil, "disk_usage") as usage, patch.object( + qualify_size, "run", return_value=f"Filesystem 1B-blocks Used Available\n/data 1 1 {provider_gib * 1024**3}" + ): + usage.return_value.free = scratch_gib * 1024**3 + qualify_size.require_size_capacity(Path("/scratch"), "provider") + + def test_accepts_both_volumes_at_threshold(self): + self.check(40, 40) + + def test_rejects_insufficient_scratch(self): + with self.assertRaisesRegex(RuntimeError, "scratch=39.0 GiB"): + self.check(39, 50) + + def test_rejects_insufficient_provider(self): + with self.assertRaisesRegex(RuntimeError, "provider=39.0 GiB"): + self.check(50, 39) + + def test_rejects_unparseable_provider_capacity(self): + with patch.object(qualify_size.shutil, "disk_usage") as usage, patch.object( + qualify_size, "run", return_value="not a df result" + ): + usage.return_value.free = 50 * 1024**3 + with self.assertRaisesRegex(RuntimeError, "could not read free space"): + qualify_size.require_size_capacity(Path("/scratch"), "provider") + + +if __name__ == "__main__": + unittest.main() From a00d474b78c3d8cec3dcbaa6e6886369c492174c Mon Sep 17 00:00:00 2001 From: forhappy Date: Tue, 29 Sep 2026 18:02:36 -0700 Subject: [PATCH 12/12] Document cold-scan movement budget bound --- docs/performance-plan.md | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/docs/performance-plan.md b/docs/performance-plan.md index 87d40b8..75edb9b 100644 --- a/docs/performance-plan.md +++ b/docs/performance-plan.md @@ -206,6 +206,15 @@ permits it. object metadata using an insertion cursor and reuses verified immutable bodies from a repository-scoped cache. Each native push/merge has a private writable generation; successful publication precedes hydration of its new objects into the shared cache. +- Demand-driven `release_idle_cell` and automatic pressure shedding share + Cellule's two-movements-per-second, two-in-flight budget. A successful + release also deletes the local SQLite state. At 100 active slots, a + first-pass scan of 10,000 distinct locally owned repositories requires at + least 9,900 releases: the budget alone implies at least 4,950 seconds + (82.5 minutes), before restore, Git work or network time. This is a lower + bound from code, not a measured scan. Bursts can instead receive a capacity + error after Canopy's single one-second retry. Do not interpret the 10,000 + identity target as a uniform-access pass until this path is qualified. - Eight shared heavy-request slots and four active and four pending slots per account prevent one account from filling every transfer wait position. The bounded Linux profile establishes containment, not latency, throughput or @@ -383,6 +392,12 @@ count dropped arrivals and timeouts as failures, never omit them from results. Run skewed and uniform distributions across 10,000 identities, each with 100, 500 and 1,000 active repositories. Vary simultaneous pack workers independently. +For uniform access, report successful releases per second, movement-budget +rejections, retry outcomes and cold-admission failures. A future demand-driven +release budget must retain generation checks, settled-work preflight and +authoritative release, while keeping automatic pressure shedding paced and +preserving bounded memory/disk use. Raise no shared movement limit solely to +make a benchmark pass. Report cold activation latency by database size, local cache state and restore bytes. Report clone/fetch first-byte latency, throughput, CPU per transferred GiB and object-store cost. Write latency includes durable acknowledgement.