Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 23 additions & 9 deletions crates/walgit-server/src/lfs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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"
Expand Down Expand Up @@ -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::<Vec<_>>();
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::<Vec<_>>()
.await;
let missing = body
.objects
.iter()
.zip(&local)
.filter(|(_, exists)| !**exists)
.map(|(o, _)| (o.oid.clone(), o.size))
.collect::<Vec<_>>();
let upstream_has = match (&cfg.upstream.lfs, missing.is_empty()) {
(Some(upstream), false) => {
st.lfs_upstream
Expand Down
58 changes: 57 additions & 1 deletion crates/walgit-server/tests/lfs_upstream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand All @@ -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,
Expand Down Expand Up @@ -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::<Vec<_>>()})
}

#[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<String> = (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::<Vec<_>>(),
});

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<()> {
Expand Down
5 changes: 5 additions & 0 deletions docs/ROUNDTRIPS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Loading