diff --git a/crates/walgit-server/src/lfs.rs b/crates/walgit-server/src/lfs.rs index da8bc30..0fd6002 100644 --- a/crates/walgit-server/src/lfs.rs +++ b/crates/walgit-server/src/lfs.rs @@ -6,6 +6,7 @@ use axum::body::Body; use axum::http::{HeaderMap, StatusCode}; use axum::response::{IntoResponse, Response}; use bytes::Bytes; +use futures::StreamExt; use serde::{Deserialize, Serialize}; use walgit_proto::keys; @@ -16,6 +17,9 @@ use crate::smart::open_repo; use crate::stream::body_to_async_read; use walgit_store::{ObjectStore, ObjectStoreExt, PutBody, PutMode}; +// Collapse git-lfs's default 100-object batch without unbounded store fan-out. +const LOCAL_PRESENCE_CONCURRENCY: usize = 16; + #[derive(Debug, Deserialize)] pub struct BatchRequest { pub operation: String, // "upload" | "download" @@ -96,15 +100,25 @@ pub async fn batch( require_lfs_oid(&o.oid)?; } // Local presence first; then one bounded upstream batch for the misses. - let mut local = Vec::with_capacity(body.objects.len()); - let mut missing = Vec::new(); - for o in &body.objects { - let exists = store.exists(&keys::lfs_key(&o.oid)).await.unwrap_or(false); - local.push(exists); - if !exists { - missing.push((o.oid.clone(), o.size)); - } - } + let local_keys = body + .objects + .iter() + .map(|o| keys::lfs_key(&o.oid)) + .collect::>(); + let local = futures::stream::iter(local_keys.into_iter().map(|key| { + let store = store.clone(); + async move { store.exists(&key).await.unwrap_or(false) } + })) + .buffered(LOCAL_PRESENCE_CONCURRENCY) + .collect::>() + .await; + let missing = body + .objects + .iter() + .zip(&local) + .filter(|(_, exists)| !**exists) + .map(|(o, _)| (o.oid.clone(), o.size)) + .collect::>(); let upstream_has = match (&cfg.upstream.lfs, missing.is_empty()) { (Some(upstream), false) => { st.lfs_upstream diff --git a/crates/walgit-server/tests/lfs_upstream.rs b/crates/walgit-server/tests/lfs_upstream.rs index 1027673..803af31 100644 --- a/crates/walgit-server/tests/lfs_upstream.rs +++ b/crates/walgit-server/tests/lfs_upstream.rs @@ -14,6 +14,7 @@ mod harness; use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; use anyhow::Result; use axum::{ @@ -26,7 +27,8 @@ use harness::Server; use serde_json::{Value, json}; use sha2::{Digest, Sha256}; use walgit_proto::keys; -use walgit_store::ObjectStoreExt; +use walgit_store::memory::MemoryStore; +use walgit_store::{ObjectStoreExt, PutMode}; struct Mock { oid: String, @@ -100,6 +102,60 @@ fn batch_req(op: &str, objects: &[(&str, usize)]) -> Value { json!({"operation": op, "transfers": ["basic"], "objects": objects.iter().map(|(o, s)| json!({"oid": o, "size": s})).collect::>()}) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn local_presence_checks_are_parallel_and_preserve_batch_order() -> Result<()> { + let mut raw_store = MemoryStore::new(); + let objects: Vec = (0..100).map(|i| format!("{i:064x}")).collect(); + for oid in objects.iter().step_by(2) { + raw_store + .put_bytes( + &format!("repos/o/r/{}", keys::lfs_key(oid)), + "present", + PutMode::Create, + ) + .await?; + } + raw_store.latency = Some(Duration::from_millis(20)); + + let server = Server::start_with_store_and_tweak(Arc::new(raw_store), |_| {}).await?; + server.put_repo("o", "r").await?; + let batch_url = format!("{}/o/r.git/info/lfs/objects/batch", server.base_url); + let request = json!({ + "operation": "download", + "transfers": ["basic"], + "objects": objects + .iter() + .map(|oid| json!({"oid": oid, "size": 7})) + .collect::>(), + }); + + let started = Instant::now(); + let response: Value = reqwest::Client::new() + .post(batch_url) + .json(&request) + .send() + .await? + .json() + .await?; + let elapsed = started.elapsed(); + + assert!( + elapsed < Duration::from_secs(1), + "100 independent presence checks took {elapsed:?}" + ); + let response_objects = response["objects"].as_array().unwrap(); + assert_eq!(response_objects.len(), objects.len()); + for (i, (actual, oid)) in response_objects.iter().zip(&objects).enumerate() { + assert_eq!(actual["oid"], oid.as_str(), "response order changed at {i}"); + if i % 2 == 0 { + assert!(actual["actions"]["download"].is_object()); + } else { + assert_eq!(actual["error"]["code"], 404); + } + } + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn upstream_objects_are_present_for_upload_and_streamed_then_persisted_for_download() -> Result<()> { diff --git a/docs/ROUNDTRIPS.md b/docs/ROUNDTRIPS.md index 08eef6e..49b704e 100644 --- a/docs/ROUNDTRIPS.md +++ b/docs/ROUNDTRIPS.md @@ -125,3 +125,8 @@ AWS documents the native conditions for [DELETE](https://docs.aws.amazon.com/Ama and [multipart completion](https://docs.aws.amazon.com/AmazonS3/latest/API/API_CompleteMultipartUpload.html). The SDK-transport tests assert headers, stale-token rejection, surviving rival data and multipart aborts. Model/negative controls and witnesses: `StoreConditions`. + +LFS batch presence checks run in ordered groups of at most 16 concurrent HEADs: +critical-path depth changes from N to ceil(N/16), with N requests unchanged. +`local_presence_checks_are_parallel_and_preserve_batch_order` covers response +ordering and presence results across 100 objects with artificial store latency.