From 31a14c77407491bda8f96124d02c58f3f25d7bfc Mon Sep 17 00:00:00 2001 From: euxaristia Date: Mon, 7 Sep 2026 02:42:56 -0400 Subject: [PATCH 1/4] fix(node): bound federated repository aggregation --- crates/gitlawb-node/src/api/repos.rs | 432 ++++++++++++++++++++++++--- 1 file changed, 388 insertions(+), 44 deletions(-) diff --git a/crates/gitlawb-node/src/api/repos.rs b/crates/gitlawb-node/src/api/repos.rs index 4e327c42..358a1461 100644 --- a/crates/gitlawb-node/src/api/repos.rs +++ b/crates/gitlawb-node/src/api/repos.rs @@ -3,6 +3,8 @@ use axum::http::StatusCode; use axum::response::Response; use axum::Json; use bytes::Bytes; +use futures::{stream, StreamExt}; +use std::future::Future; use std::sync::Arc; use crate::auth::{caller_authorized_to_push, AuthenticatedDid}; @@ -21,6 +23,16 @@ use crate::webhooks; /// The git all-zeros object id — the create/delete sentinel in a ref update. const ZERO_SHA: &str = "0000000000000000000000000000000000000000"; +const FEDERATED_PEER_CONCURRENCY: usize = 4; +const FEDERATED_PEER_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); +const MAX_FEDERATED_PEER_BYTES: usize = 512 * 1024; +const MAX_FEDERATED_PEER_REPOS: usize = 200; +const MAX_FEDERATED_REPOS: usize = 1_000; +const MAX_FEDERATED_REPO_JSON_BYTES: usize = 2 * 1024 * 1024; +// The aggregate byte budget counts serialized repo values and separators. Keep +// a fixed reserve for the response object's keys, counters, and closing bytes. +const FEDERATED_RESPONSE_OVERHEAD_BYTES: usize = 256; + /// The set of blob OIDs withheld from **anonymous** replication for a repo, or /// `None` when the repo must not replicate at all (private / mode A / /// undetermined — fail closed). This is the anonymous replication gate: @@ -2897,6 +2909,172 @@ pub async fn list_refs( )) } +struct FederatedPeerRepos { + node_url: String, + node_did: String, + repos: Vec, +} + +struct BoundedFederatedPeerRows(Vec); + +impl<'de> serde::Deserialize<'de> for BoundedFederatedPeerRows { + fn deserialize(deserializer: D) -> std::result::Result + where + D: serde::Deserializer<'de>, + { + struct RowsVisitor; + + impl<'de> serde::de::Visitor<'de> for RowsVisitor { + type Value = BoundedFederatedPeerRows; + + fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str("a bounded array of repository objects") + } + + fn visit_seq(self, mut rows: A) -> std::result::Result + where + A: serde::de::SeqAccess<'de>, + { + let mut repos = Vec::with_capacity(MAX_FEDERATED_PEER_REPOS); + while let Some(repo) = rows.next_element::()? { + if repos.len() == MAX_FEDERATED_PEER_REPOS { + return Err(::custom( + "peer repository row limit exceeded", + )); + } + if !repo.is_object() { + return Err(::custom( + "peer repository row is not an object", + )); + } + repos.push(repo); + } + Ok(BoundedFederatedPeerRows(repos)) + } + } + + deserializer.deserialize_seq(RowsVisitor) + } +} + +struct FederatedRepoBudget { + repos: Vec, + json_bytes: usize, + max_repos: usize, + max_json_bytes: usize, + truncated: bool, +} + +impl FederatedRepoBudget { + fn new(capacity: usize, max_repos: usize, max_json_bytes: usize) -> Self { + Self { + repos: Vec::with_capacity(capacity.min(max_repos)), + json_bytes: FEDERATED_RESPONSE_OVERHEAD_BYTES, + max_repos, + max_json_bytes, + truncated: false, + } + } + + fn push(&mut self, repo: serde_json::Value) -> bool { + if self.is_saturated() { + self.truncated = true; + return false; + } + let Ok(repo_bytes) = serde_json::to_vec(&repo) else { + return true; + }; + let charged = repo_bytes.len().saturating_add(1); + if self.json_bytes.saturating_add(charged) > self.max_json_bytes { + self.truncated = true; + return false; + } + self.json_bytes += charged; + self.repos.push(repo); + true + } + + fn is_saturated(&self) -> bool { + self.repos.len() >= self.max_repos || self.json_bytes >= self.max_json_bytes + } +} + +async fn fetch_federated_peer_repos( + client: &reqwest::Client, + url: &str, +) -> Option> { + let mut response = client.get(url).send().await.ok()?; + if !response.status().is_success() + || response + .content_length() + .is_some_and(|len| len > MAX_FEDERATED_PEER_BYTES as u64) + { + return None; + } + + let capacity = response + .content_length() + .and_then(|len| usize::try_from(len).ok()) + .unwrap_or(0) + .min(MAX_FEDERATED_PEER_BYTES); + let mut body = Vec::with_capacity(capacity); + while let Some(chunk) = response.chunk().await.ok()? { + if body.len().saturating_add(chunk.len()) > MAX_FEDERATED_PEER_BYTES { + return None; + } + body.extend_from_slice(&chunk); + } + + let BoundedFederatedPeerRows(repos) = serde_json::from_slice(&body).ok()?; + Some(repos) +} + +fn enrich_federated_repo( + mut repo: serde_json::Value, + node_url: &str, + node_did: &str, +) -> Option { + let object = repo.as_object_mut()?; + object.insert( + "node_url".to_string(), + serde_json::Value::String(node_url.to_string()), + ); + object.insert( + "node_did".to_string(), + serde_json::Value::String(node_did.to_string()), + ); + object.insert("local".to_string(), serde_json::Value::Bool(false)); + Some(repo) +} + +async fn collect_federated_fetches(fetches: I, aggregate: &mut FederatedRepoBudget) -> usize +where + I: IntoIterator, + F: Future>, +{ + let mut responses = stream::iter(fetches).buffer_unordered(FEDERATED_PEER_CONCURRENCY); + let mut nodes_queried = 0; + + while let Some(response) = responses.next().await { + let Some(response) = response else { + continue; + }; + nodes_queried += 1; + for repo in response.repos { + let Some(repo) = enrich_federated_repo(repo, &response.node_url, &response.node_did) + else { + continue; + }; + if !aggregate.push(repo) || aggregate.is_saturated() { + aggregate.truncated = true; + return nodes_queried; + } + } + } + + nodes_queried +} + /// GET /api/v1/repos/federated /// /// Query all known peers for their public repos and return a merged view of @@ -2929,66 +3107,69 @@ pub async fn list_federated_repos( .unwrap_or_else(|| "http://127.0.0.1:7545".to_string()); let local_node_did = state.node_did.to_string(); - let mut all_repos: Vec = Vec::with_capacity(local_repos.len()); + let mut aggregate = FederatedRepoBudget::new( + local_repos.len(), + MAX_FEDERATED_REPOS, + MAX_FEDERATED_REPO_JSON_BYTES, + ); for (r, count) in &local_repos { - let mut v = serde_json::to_value(to_response(r, &state, *count)).unwrap_or_default(); - v["node_url"] = serde_json::Value::String(local_node_url.clone()); - v["node_did"] = serde_json::Value::String(local_node_did.clone()); - v["local"] = serde_json::Value::Bool(true); - all_repos.push(v); + let Ok(mut repo) = serde_json::to_value(to_response(r, &state, *count)) else { + continue; + }; + repo["node_url"] = serde_json::Value::String(local_node_url.clone()); + repo["node_did"] = serde_json::Value::String(local_node_did.clone()); + repo["local"] = serde_json::Value::Bool(true); + if !aggregate.push(repo) || aggregate.is_saturated() { + aggregate.truncated = true; + break; + } } - // Query peers in parallel + // Keep peer work bounded even when the table contains many reachable rows. + // Dropping the buffered stream at the aggregate ceiling cancels its in-flight + // request futures and leaves the remaining iterator entries unpolled. let peers = state.db.list_peers().await.unwrap_or_default(); - let client = &state.http_client; - - let fetch_tasks: Vec<_> = peers + let fetches: Vec<_> = peers .into_iter() .filter(|p| p.last_ping_ok && !p.http_url.is_empty()) .map(|peer| { - let client = Arc::clone(client); - let url = format!("{}/api/v1/repos", peer.http_url.trim_end_matches('/')); + let client = Arc::clone(&state.http_client); + let url = format!( + "{}/api/v1/repos?limit={MAX_FEDERATED_PEER_REPOS}&offset=0", + peer.http_url.trim_end_matches('/') + ); let peer_did = peer.did.clone(); let peer_url = peer.http_url.clone(); - tokio::spawn(async move { - let result = tokio::time::timeout( - std::time::Duration::from_secs(5), - client.get(&url).send(), + async move { + tokio::time::timeout( + FEDERATED_PEER_TIMEOUT, + fetch_federated_peer_repos(&client, &url), ) - .await; - match result { - Ok(Ok(resp)) if resp.status().is_success() => { - if let Ok(repos) = resp.json::>().await { - let enriched: Vec = repos - .into_iter() - .map(|mut r| { - r["node_url"] = serde_json::Value::String(peer_url.clone()); - r["node_did"] = serde_json::Value::String(peer_did.clone()); - r["local"] = serde_json::Value::Bool(false); - r - }) - .collect(); - return enriched; - } - } - _ => {} - } - vec![] - }) + .await + .ok() + .flatten() + .map(|repos| FederatedPeerRepos { + node_url: peer_url, + node_did: peer_did, + repos, + }) + } }) .collect(); - for task in fetch_tasks { - if let Ok(repos) = task.await { - all_repos.extend(repos); - } - } + let peer_nodes_queried = if aggregate.is_saturated() { + 0 + } else { + collect_federated_fetches(fetches, &mut aggregate).await + }; - let count = all_repos.len(); + let count = aggregate.repos.len(); + let truncated = aggregate.truncated; Ok(Json(serde_json::json!({ - "repos": all_repos, + "repos": aggregate.repos, "count": count, - "nodes_queried": 1, // local + peers that responded + "nodes_queried": 1 + peer_nodes_queried, + "truncated": truncated, }))) } @@ -3343,6 +3524,169 @@ mod tests { const OWNER_SHORT: &str = "z6MkpTHR8VNsBxYAAWHut2Geadd9jSwuBV8xRoAnwWsdvktH"; const STRANGER_DID: &str = "did:key:z6Mkffonly5tranger0000000000000000000000000000000"; + #[tokio::test] + async fn federated_peer_response_enforces_byte_and_row_ceilings() { + let mut row_server = mockito::Server::new_async().await; + let row_body = + serde_json::to_string(&vec![serde_json::json!({}); MAX_FEDERATED_PEER_REPOS + 1]) + .unwrap(); + let row_mock = row_server + .mock("GET", "/api/v1/repos") + .with_status(200) + .with_header("content-type", "application/json") + .with_body(row_body) + .expect(1) + .create_async() + .await; + + let client = reqwest::Client::new(); + assert!( + fetch_federated_peer_repos(&client, &format!("{}/api/v1/repos", row_server.url())) + .await + .is_none(), + "a peer page above the row ceiling must be omitted" + ); + row_mock.assert_async().await; + + let byte_body = serde_json::to_string(&vec![serde_json::json!({ + "description": "x".repeat(MAX_FEDERATED_PEER_BYTES) + })]) + .unwrap(); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let byte_server = tokio::spawn(async move { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let (mut socket, _) = listener.accept().await.unwrap(); + let mut request = [0; 1_024]; + let _ = socket.read(&mut request).await.unwrap(); + socket + .write_all( + format!( + "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\n\ + transfer-encoding: chunked\r\nconnection: close\r\n\r\n{:x}\r\n", + byte_body.len() + ) + .as_bytes(), + ) + .await + .unwrap(); + let _ = socket.write_all(byte_body.as_bytes()).await; + let _ = socket.write_all(b"\r\n0\r\n\r\n").await; + }); + + assert!( + fetch_federated_peer_repos(&client, &format!("http://{address}/api/v1/repos")) + .await + .is_none(), + "a chunked peer body above the byte ceiling must be omitted without Content-Length" + ); + byte_server.await.unwrap(); + } + + #[test] + fn federated_aggregate_stays_inside_row_and_serialized_byte_budgets() { + let mut rows = FederatedRepoBudget::new(3, 2, 4_096); + assert!(rows.push(serde_json::json!({ "name": "one" }))); + assert!(rows.push(serde_json::json!({ "name": "two" }))); + assert!(rows.is_saturated()); + assert!(!rows.push(serde_json::json!({ "name": "three" }))); + assert_eq!(rows.repos.len(), 2); + assert!(rows.truncated); + + let max_bytes = 600; + let mut bytes = FederatedRepoBudget::new(8, 8, max_bytes); + while bytes.push(serde_json::json!({ "description": "x".repeat(100) })) {} + let count = bytes.repos.len(); + let truncated = bytes.truncated; + let response = serde_json::json!({ + "repos": bytes.repos, + "count": count, + "nodes_queried": 1, + "truncated": truncated, + }); + assert!(serde_json::to_vec(&response).unwrap().len() <= max_bytes); + assert!(response["truncated"].as_bool().unwrap()); + } + + struct ActiveFetch { + active: std::sync::Arc, + } + + impl Drop for ActiveFetch { + fn drop(&mut self) { + self.active + .fetch_sub(1, std::sync::atomic::Ordering::SeqCst); + } + } + + #[tokio::test] + async fn federated_collection_bounds_concurrency_and_cancels_at_aggregate_limit() { + use std::sync::atomic::{AtomicUsize, Ordering}; + + let active = Arc::new(AtomicUsize::new(0)); + let max_active = Arc::new(AtomicUsize::new(0)); + let started = Arc::new(AtomicUsize::new(0)); + let mut senders = Vec::new(); + let mut fetches = Vec::new(); + + for i in 0..10 { + let (tx, rx) = tokio::sync::oneshot::channel(); + senders.push(Some(tx)); + let active = Arc::clone(&active); + let max_active = Arc::clone(&max_active); + let started = Arc::clone(&started); + fetches.push(async move { + started.fetch_add(1, Ordering::SeqCst); + let current = active.fetch_add(1, Ordering::SeqCst) + 1; + max_active.fetch_max(current, Ordering::SeqCst); + let _active = ActiveFetch { active }; + rx.await.ok()?; + Some(FederatedPeerRepos { + node_url: format!("https://peer-{i}.example"), + node_did: format!("did:key:peer-{i}"), + repos: vec![serde_json::json!({ "id": i })], + }) + }); + } + + let task = tokio::spawn(async move { + let mut aggregate = FederatedRepoBudget::new(1, 1, usize::MAX); + let nodes = collect_federated_fetches(fetches, &mut aggregate).await; + (aggregate, nodes) + }); + tokio::time::timeout(std::time::Duration::from_secs(1), async { + while started.load(Ordering::SeqCst) < FEDERATED_PEER_CONCURRENCY { + tokio::task::yield_now().await; + } + }) + .await + .expect("the bounded fetch window should fill"); + + assert_eq!(started.load(Ordering::SeqCst), FEDERATED_PEER_CONCURRENCY); + assert_eq!(active.load(Ordering::SeqCst), FEDERATED_PEER_CONCURRENCY); + assert_eq!( + max_active.load(Ordering::SeqCst), + FEDERATED_PEER_CONCURRENCY + ); + + senders[0].take().unwrap().send(()).unwrap(); + let (aggregate, nodes) = tokio::time::timeout(std::time::Duration::from_secs(1), task) + .await + .expect("the aggregate ceiling should stop collection") + .unwrap(); + + assert_eq!(nodes, 1); + assert_eq!(aggregate.repos.len(), 1); + assert!(aggregate.truncated); + assert_eq!(started.load(Ordering::SeqCst), FEDERATED_PEER_CONCURRENCY); + assert_eq!( + active.load(Ordering::SeqCst), + 0, + "in-flight fetches must be dropped" + ); + } + #[test] fn upload_pack_request_finalizes_only_with_done_pktline() { let want = "0032want 1111111111111111111111111111111111111111\n"; From 91dd5d108edd964bcd9c15236f2866c5f184d7cf Mon Sep 17 00:00:00 2001 From: euxaristia Date: Mon, 7 Sep 2026 02:50:27 -0400 Subject: [PATCH 2/4] fix(node): propagate federated peer truncation --- crates/gitlawb-node/src/api/repos.rs | 91 ++++++++++++++++++++++++++-- 1 file changed, 87 insertions(+), 4 deletions(-) diff --git a/crates/gitlawb-node/src/api/repos.rs b/crates/gitlawb-node/src/api/repos.rs index 358a1461..69c11446 100644 --- a/crates/gitlawb-node/src/api/repos.rs +++ b/crates/gitlawb-node/src/api/repos.rs @@ -2913,6 +2913,12 @@ struct FederatedPeerRepos { node_url: String, node_did: String, repos: Vec, + truncated: bool, +} + +struct FederatedPeerPage { + repos: Vec, + truncated: bool, } struct BoundedFederatedPeerRows(Vec); @@ -3002,7 +3008,7 @@ impl FederatedRepoBudget { async fn fetch_federated_peer_repos( client: &reqwest::Client, url: &str, -) -> Option> { +) -> Option { let mut response = client.get(url).send().await.ok()?; if !response.status().is_success() || response @@ -3011,6 +3017,11 @@ async fn fetch_federated_peer_repos( { return None; } + let total = response + .headers() + .get("x-total-count") + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.trim().parse::().ok()); let capacity = response .content_length() @@ -3026,7 +3037,8 @@ async fn fetch_federated_peer_repos( } let BoundedFederatedPeerRows(repos) = serde_json::from_slice(&body).ok()?; - Some(repos) + let truncated = total.is_some_and(|total| total > repos.len() as u64); + Some(FederatedPeerPage { repos, truncated }) } fn enrich_federated_repo( @@ -3060,6 +3072,7 @@ where continue; }; nodes_queried += 1; + aggregate.truncated |= response.truncated; for repo in response.repos { let Some(repo) = enrich_federated_repo(repo, &response.node_url, &response.node_did) else { @@ -3148,10 +3161,11 @@ pub async fn list_federated_repos( .await .ok() .flatten() - .map(|repos| FederatedPeerRepos { + .map(|page| FederatedPeerRepos { node_url: peer_url, node_did: peer_did, - repos, + repos: page.repos, + truncated: page.truncated, }) } }) @@ -3584,6 +3598,74 @@ mod tests { byte_server.await.unwrap(); } + #[tokio::test] + async fn federated_peer_total_count_prevents_false_complete_page() { + let mut server = mockito::Server::new_async().await; + let body = serde_json::to_string( + &(0..MAX_FEDERATED_PEER_REPOS) + .map(|id| serde_json::json!({ "id": id })) + .collect::>(), + ) + .unwrap(); + let page_mock = server + .mock("GET", "/api/v1/repos") + .with_status(200) + .with_header("content-type", "application/json") + .with_header("x-total-count", &(MAX_FEDERATED_PEER_REPOS + 1).to_string()) + .with_body(body) + .expect(1) + .create_async() + .await; + + let page = fetch_federated_peer_repos( + &reqwest::Client::new(), + &format!("{}/api/v1/repos", server.url()), + ) + .await + .expect("the capped peer page is otherwise valid"); + page_mock.assert_async().await; + assert_eq!(page.repos.len(), MAX_FEDERATED_PEER_REPOS); + assert!(page.truncated, "the peer total proves another row exists"); + + let mut aggregate = FederatedRepoBudget::new(300, 300, MAX_FEDERATED_REPO_JSON_BYTES); + let nodes = collect_federated_fetches( + [std::future::ready(Some(FederatedPeerRepos { + node_url: server.url(), + node_did: "did:key:peer".to_string(), + repos: page.repos, + truncated: page.truncated, + }))], + &mut aggregate, + ) + .await; + assert_eq!(nodes, 1); + assert_eq!(aggregate.repos.len(), MAX_FEDERATED_PEER_REPOS); + assert!( + aggregate.truncated, + "the aggregate must not report the paged peer as complete" + ); + + let mut malformed_server = mockito::Server::new_async().await; + let malformed_mock = malformed_server + .mock("GET", "/api/v1/repos") + .with_status(200) + .with_header("content-type", "application/json") + .with_header("x-total-count", "not-a-count") + .with_body("[{}]") + .expect(1) + .create_async() + .await; + let malformed = fetch_federated_peer_repos( + &reqwest::Client::new(), + &format!("{}/api/v1/repos", malformed_server.url()), + ) + .await + .expect("a malformed optional total must not discard a bounded page"); + malformed_mock.assert_async().await; + assert!(!malformed.truncated); + assert_eq!(malformed.repos.len(), 1); + } + #[test] fn federated_aggregate_stays_inside_row_and_serialized_byte_budgets() { let mut rows = FederatedRepoBudget::new(3, 2, 4_096); @@ -3646,6 +3728,7 @@ mod tests { node_url: format!("https://peer-{i}.example"), node_did: format!("did:key:peer-{i}"), repos: vec![serde_json::json!({ "id": i })], + truncated: false, }) }); } From 5b62d440ce0e51b71d864295a071934011802a30 Mon Sep 17 00:00:00 2001 From: euxaristia Date: Tue, 8 Sep 2026 20:32:55 -0400 Subject: [PATCH 3/4] fix(node): Bound federation latency and preserve truncated peer pages. --- crates/gitlawb-node/src/api/repos.rs | 104 ++++++++++++++++++++++----- crates/gitlawb-node/src/db/mod.rs | 20 ++++++ crates/gitlawb-node/src/server.rs | 14 +++- 3 files changed, 120 insertions(+), 18 deletions(-) diff --git a/crates/gitlawb-node/src/api/repos.rs b/crates/gitlawb-node/src/api/repos.rs index 69c11446..45ab7579 100644 --- a/crates/gitlawb-node/src/api/repos.rs +++ b/crates/gitlawb-node/src/api/repos.rs @@ -2921,7 +2921,7 @@ struct FederatedPeerPage { truncated: bool, } -struct BoundedFederatedPeerRows(Vec); +struct BoundedFederatedPeerRows(Vec, bool); impl<'de> serde::Deserialize<'de> for BoundedFederatedPeerRows { fn deserialize(deserializer: D) -> std::result::Result @@ -2942,12 +2942,10 @@ impl<'de> serde::Deserialize<'de> for BoundedFederatedPeerRows { A: serde::de::SeqAccess<'de>, { let mut repos = Vec::with_capacity(MAX_FEDERATED_PEER_REPOS); - while let Some(repo) = rows.next_element::()? { - if repos.len() == MAX_FEDERATED_PEER_REPOS { - return Err(::custom( - "peer repository row limit exceeded", - )); - } + while repos.len() < MAX_FEDERATED_PEER_REPOS { + let Some(repo) = rows.next_element::()? else { + return Ok(BoundedFederatedPeerRows(repos, false)); + }; if !repo.is_object() { return Err(::custom( "peer repository row is not an object", @@ -2955,7 +2953,11 @@ impl<'de> serde::Deserialize<'de> for BoundedFederatedPeerRows { } repos.push(repo); } - Ok(BoundedFederatedPeerRows(repos)) + let mut truncated = false; + while rows.next_element::()?.is_some() { + truncated = true; + } + Ok(BoundedFederatedPeerRows(repos, truncated)) } } @@ -3036,8 +3038,8 @@ async fn fetch_federated_peer_repos( body.extend_from_slice(&chunk); } - let BoundedFederatedPeerRows(repos) = serde_json::from_slice(&body).ok()?; - let truncated = total.is_some_and(|total| total > repos.len() as u64); + let BoundedFederatedPeerRows(repos, overflow) = serde_json::from_slice(&body).ok()?; + let truncated = overflow || total.is_some_and(|total| total > repos.len() as u64); Some(FederatedPeerPage { repos, truncated }) } @@ -3067,8 +3069,20 @@ where let mut responses = stream::iter(fetches).buffer_unordered(FEDERATED_PEER_CONCURRENCY); let mut nodes_queried = 0; - while let Some(response) = responses.next().await { + let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(10); + loop { + let response = match tokio::time::timeout_at(deadline, responses.next()).await { + Ok(Some(response)) => response, + Ok(None) => break, + Err(_) => { + aggregate.truncated = true; + tracing::warn!("federated aggregation deadline exceeded"); + break; + } + }; let Some(response) = response else { + aggregate.truncated = true; + tracing::warn!("federated peer fetch failed or exceeded its response limits"); continue; }; nodes_queried += 1; @@ -3092,7 +3106,8 @@ where /// /// Query all known peers for their public repos and return a merged view of /// the network. Each repo includes a `node_url` and `node_did` indicating -/// which node hosts it. Results from unreachable peers are silently omitted. +/// which node hosts it. Peer work stops after ten seconds or 200 peers; omitted +/// results set `truncated`. The route admits 12 requests per minute per client IP. pub async fn list_federated_repos( State(state): State, auth: Option>, @@ -3141,7 +3156,10 @@ pub async fn list_federated_repos( // Keep peer work bounded even when the table contains many reachable rows. // Dropping the buffered stream at the aggregate ceiling cancels its in-flight // request futures and leaves the remaining iterator entries unpolled. - let peers = state.db.list_peers().await.unwrap_or_default(); + // One extra row detects omitted peers without materializing the full table. + let mut peers = state.db.list_federation_peers(201).await?; + aggregate.truncated |= peers.len() > 200; + peers.truncate(200); let fetches: Vec<_> = peers .into_iter() .filter(|p| p.last_ping_ok && !p.http_url.is_empty()) @@ -3538,6 +3556,58 @@ mod tests { const OWNER_SHORT: &str = "z6MkpTHR8VNsBxYAAWHut2Geadd9jSwuBV8xRoAnwWsdvktH"; const STRANGER_DID: &str = "did:key:z6Mkffonly5tranger0000000000000000000000000000000"; + #[sqlx::test] + async fn federated_peer_query_is_bounded(pool: sqlx::PgPool) { + let state = crate::test_support::test_state(pool.clone()).await; + sqlx::query("INSERT INTO peers (did, http_url, last_ping_ok, announced_at) SELECT 'peer-' || n, 'https://example.com', TRUE, '2026-01-01T00:00:00Z' FROM generate_series(1, 205) n") + .execute(&pool).await.unwrap(); + let peers = state.db.list_federation_peers(1000).await.unwrap(); + assert_eq!(peers.len(), 201); + assert_eq!(state.db.list_federation_peers(2).await.unwrap().len(), 2); + assert!(state.db.list_federation_peers(0).await.unwrap().is_empty()); + } + + #[sqlx::test] + async fn federated_route_brakes_repeated_anonymous_requests(pool: sqlx::PgPool) { + use tower::ServiceExt; + let mut state = crate::test_support::test_state(pool).await; + state.push_limiter_trust = crate::rate_limit::TrustedProxy::None; + let router = crate::server::build_router(state); + for n in 0..13 { + let mut request = axum::http::Request::builder() + .uri("/api/v1/repos/federated") + .body(axum::body::Body::empty()) + .unwrap(); + request.extensions_mut().insert(axum::extract::ConnectInfo( + "203.0.113.42:5000".parse::().unwrap(), + )); + let response = router.clone().oneshot(request).await.unwrap(); + assert_eq!( + response.status(), + if n < 12 { + StatusCode::OK + } else { + StatusCode::TOO_MANY_REQUESTS + } + ); + } + } + + #[tokio::test] + async fn federated_aggregate_deadline_returns_partial_results() { + let mut budget = FederatedRepoBudget::new(0, 1000, MAX_FEDERATED_REPO_JSON_BYTES); + let fetches = + (0..8).map(|_| async { std::future::pending::>().await }); + let nodes = tokio::time::timeout( + std::time::Duration::from_secs(12), + collect_federated_fetches(fetches, &mut budget), + ) + .await + .expect("the aggregate must stop before the caller deadline"); + assert_eq!(nodes, 0); + assert!(budget.truncated); + } + #[tokio::test] async fn federated_peer_response_enforces_byte_and_row_ceilings() { let mut row_server = mockito::Server::new_async().await; @@ -3554,12 +3624,12 @@ mod tests { .await; let client = reqwest::Client::new(); - assert!( + let page = fetch_federated_peer_repos(&client, &format!("{}/api/v1/repos", row_server.url())) .await - .is_none(), - "a peer page above the row ceiling must be omitted" - ); + .unwrap(); + assert_eq!(page.repos.len(), MAX_FEDERATED_PEER_REPOS); + assert!(page.truncated); row_mock.assert_async().await; let byte_body = serde_json::to_string(&vec![serde_json::json!({ diff --git a/crates/gitlawb-node/src/db/mod.rs b/crates/gitlawb-node/src/db/mod.rs index cc2cf0bd..a6eb9f42 100644 --- a/crates/gitlawb-node/src/db/mod.rs +++ b/crates/gitlawb-node/src/db/mod.rs @@ -2709,6 +2709,26 @@ impl Db { Ok(()) } + pub async fn list_federation_peers(&self, limit: i64) -> Result> { + let rows = sqlx::query( + "SELECT did, http_url, last_seen, last_ping_ok, announced_at + FROM peers WHERE last_ping_ok = TRUE AND http_url <> '' ORDER BY last_seen DESC NULLS LAST, did LIMIT $1", + ) + .bind(limit.clamp(0, 201)) + .fetch_all(&self.pool) + .await?; + Ok(rows + .into_iter() + .map(|r| PeerRecord { + did: r.get("did"), + http_url: r.get("http_url"), + last_seen: r.get("last_seen"), + last_ping_ok: r.get::("last_ping_ok"), + announced_at: r.get("announced_at"), + }) + .collect()) + } + pub async fn list_peers(&self) -> Result> { let rows = sqlx::query( "SELECT did, http_url, last_seen, last_ping_ok, announced_at diff --git a/crates/gitlawb-node/src/server.rs b/crates/gitlawb-node/src/server.rs index de61fcbe..cd38e55b 100644 --- a/crates/gitlawb-node/src/server.rs +++ b/crates/gitlawb-node/src/server.rs @@ -341,7 +341,19 @@ pub fn build_router(state: AppState) -> Router { // ── Read routes — open for public repos ─────────────────────────────── let read_routes = Router::new() .route("/api/v1/repos", get(repos::list_repos)) - .route("/api/v1/repos/federated", get(repos::list_federated_repos)) + .route( + "/api/v1/repos/federated", + get(repos::list_federated_repos) + .route_layer(middleware::from_fn(rate_limit::rate_limit_by_ip)) + .route_layer(axum::Extension(rate_limit::IpRateLimiter { + limiter: rate_limit::RateLimiter::new_bounded( + 12, + std::time::Duration::from_secs(60), + 10_000, + ), + trust: state.push_limiter_trust, + })), + ) .route("/api/v1/repos/{owner}/{repo}", get(repos::get_repo)) .route( "/api/v1/repos/{owner}/{repo}/commits", From deee4ad27a55e95d3ba9befd80c9cc205a7ba041 Mon Sep 17 00:00:00 2001 From: euxaristia Date: Wed, 9 Sep 2026 01:29:06 -0400 Subject: [PATCH 4/4] fix(node): complete federation review follow-ups Register the bounded peer-query fixture in the writer ledger, name federation bounds, and share the state-owned limiter across routers and the cleanup sweep. Cover exhausted-bucket rejection before database access. --- crates/gitlawb-node/src/api/repos.rs | 60 ++++++++++++++++++++----- crates/gitlawb-node/src/auth/mod.rs | 1 + crates/gitlawb-node/src/db/mod.rs | 4 +- crates/gitlawb-node/src/main.rs | 5 +++ crates/gitlawb-node/src/server.rs | 6 +-- crates/gitlawb-node/src/state.rs | 4 ++ crates/gitlawb-node/src/test_support.rs | 1 + 7 files changed, 64 insertions(+), 17 deletions(-) diff --git a/crates/gitlawb-node/src/api/repos.rs b/crates/gitlawb-node/src/api/repos.rs index 45ab7579..600ed5fa 100644 --- a/crates/gitlawb-node/src/api/repos.rs +++ b/crates/gitlawb-node/src/api/repos.rs @@ -25,6 +25,8 @@ const ZERO_SHA: &str = "0000000000000000000000000000000000000000"; const FEDERATED_PEER_CONCURRENCY: usize = 4; const FEDERATED_PEER_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); +const FEDERATED_AGGREGATE_DEADLINE: std::time::Duration = std::time::Duration::from_secs(10); +pub(crate) const MAX_FEDERATED_PEERS: usize = 200; const MAX_FEDERATED_PEER_BYTES: usize = 512 * 1024; const MAX_FEDERATED_PEER_REPOS: usize = 200; const MAX_FEDERATED_REPOS: usize = 1_000; @@ -3069,7 +3071,7 @@ where let mut responses = stream::iter(fetches).buffer_unordered(FEDERATED_PEER_CONCURRENCY); let mut nodes_queried = 0; - let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(10); + let deadline = tokio::time::Instant::now() + FEDERATED_AGGREGATE_DEADLINE; loop { let response = match tokio::time::timeout_at(deadline, responses.next()).await { Ok(Some(response)) => response, @@ -3106,7 +3108,8 @@ where /// /// Query all known peers for their public repos and return a merged view of /// the network. Each repo includes a `node_url` and `node_did` indicating -/// which node hosts it. Peer work stops after ten seconds or 200 peers; omitted +/// which node hosts it. Peer work is bounded by `FEDERATED_AGGREGATE_DEADLINE` +/// and `MAX_FEDERATED_PEERS`; omitted /// results set `truncated`. The route admits 12 requests per minute per client IP. pub async fn list_federated_repos( State(state): State, @@ -3157,9 +3160,12 @@ pub async fn list_federated_repos( // Dropping the buffered stream at the aggregate ceiling cancels its in-flight // request futures and leaves the remaining iterator entries unpolled. // One extra row detects omitted peers without materializing the full table. - let mut peers = state.db.list_federation_peers(201).await?; - aggregate.truncated |= peers.len() > 200; - peers.truncate(200); + let mut peers = state + .db + .list_federation_peers(MAX_FEDERATED_PEERS as i64 + 1) + .await?; + aggregate.truncated |= peers.len() > MAX_FEDERATED_PEERS; + peers.truncate(MAX_FEDERATED_PEERS); let fetches: Vec<_> = peers .into_iter() .filter(|p| p.last_ping_ok && !p.http_url.is_empty()) @@ -3559,10 +3565,11 @@ mod tests { #[sqlx::test] async fn federated_peer_query_is_bounded(pool: sqlx::PgPool) { let state = crate::test_support::test_state(pool.clone()).await; - sqlx::query("INSERT INTO peers (did, http_url, last_ping_ok, announced_at) SELECT 'peer-' || n, 'https://example.com', TRUE, '2026-01-01T00:00:00Z' FROM generate_series(1, 205) n") + sqlx::query("INSERT INTO peers (did, http_url, last_ping_ok, announced_at) SELECT 'peer-' || n, 'https://example.com', TRUE, '2026-01-01T00:00:00Z' FROM generate_series(1, $1) n") + .bind(MAX_FEDERATED_PEERS as i64 + 5) .execute(&pool).await.unwrap(); - let peers = state.db.list_federation_peers(1000).await.unwrap(); - assert_eq!(peers.len(), 201); + let peers = state.db.list_federation_peers(i64::MAX).await.unwrap(); + assert_eq!(peers.len(), MAX_FEDERATED_PEERS + 1); assert_eq!(state.db.list_federation_peers(2).await.unwrap().len(), 2); assert!(state.db.list_federation_peers(0).await.unwrap().is_empty()); } @@ -3572,8 +3579,10 @@ mod tests { use tower::ServiceExt; let mut state = crate::test_support::test_state(pool).await; state.push_limiter_trust = crate::rate_limit::TrustedProxy::None; + state.federated_rate_limiter = + crate::rate_limit::RateLimiter::new_bounded(1, std::time::Duration::from_secs(60), 10); let router = crate::server::build_router(state); - for n in 0..13 { + for n in 0..2 { let mut request = axum::http::Request::builder() .uri("/api/v1/repos/federated") .body(axum::body::Body::empty()) @@ -3584,7 +3593,7 @@ mod tests { let response = router.clone().oneshot(request).await.unwrap(); assert_eq!( response.status(), - if n < 12 { + if n == 0 { StatusCode::OK } else { StatusCode::TOO_MANY_REQUESTS @@ -3593,13 +3602,42 @@ mod tests { } } + #[tokio::test] + async fn federated_route_uses_the_state_owned_bucket_before_database_access() { + use tower::ServiceExt; + let mut state = crate::test_support::test_state_lazy(); + state.push_limiter_trust = crate::rate_limit::TrustedProxy::None; + state.federated_rate_limiter = + crate::rate_limit::RateLimiter::new_bounded(1, std::time::Duration::from_secs(60), 10); + assert!(state.federated_rate_limiter.check("203.0.113.42").await); + // Both routers must use the same exhausted bucket from AppState. + for router in [ + crate::server::build_router(state.clone()), + crate::server::build_router(state), + ] { + let mut request = axum::http::Request::builder() + .uri("/api/v1/repos/federated") + .body(axum::body::Body::empty()) + .unwrap(); + request.extensions_mut().insert(axum::extract::ConnectInfo( + "203.0.113.42:5000".parse::().unwrap(), + )); + let response = + tokio::time::timeout(std::time::Duration::from_secs(1), router.oneshot(request)) + .await + .expect("the exhausted bucket must reject before accessing the lazy database") + .unwrap(); + assert_eq!(response.status(), StatusCode::TOO_MANY_REQUESTS); + } + } + #[tokio::test] async fn federated_aggregate_deadline_returns_partial_results() { let mut budget = FederatedRepoBudget::new(0, 1000, MAX_FEDERATED_REPO_JSON_BYTES); let fetches = (0..8).map(|_| async { std::future::pending::>().await }); let nodes = tokio::time::timeout( - std::time::Duration::from_secs(12), + FEDERATED_AGGREGATE_DEADLINE + std::time::Duration::from_secs(2), collect_federated_fetches(fetches, &mut budget), ) .await diff --git a/crates/gitlawb-node/src/auth/mod.rs b/crates/gitlawb-node/src/auth/mod.rs index 27b67786..317c2094 100644 --- a/crates/gitlawb-node/src/auth/mod.rs +++ b/crates/gitlawb-node/src/auth/mod.rs @@ -530,6 +530,7 @@ mod tests { push_limiter_trust: crate::rate_limit::TrustedProxy::None, sync_trigger_rate_limiter: RateLimiter::new(60, Duration::from_secs(3600)), peer_write_rate_limiter: RateLimiter::new(600, Duration::from_secs(3600)), + federated_rate_limiter: RateLimiter::new_bounded(12, Duration::from_secs(60), 10_000), shutdown_tx: tokio::sync::watch::channel(false).0, git_read_semaphore: Arc::new(tokio::sync::Semaphore::new(64)), git_write_semaphore: Arc::new(tokio::sync::Semaphore::new(64)), diff --git a/crates/gitlawb-node/src/db/mod.rs b/crates/gitlawb-node/src/db/mod.rs index a6eb9f42..8892c5a8 100644 --- a/crates/gitlawb-node/src/db/mod.rs +++ b/crates/gitlawb-node/src/db/mod.rs @@ -2714,7 +2714,7 @@ impl Db { "SELECT did, http_url, last_seen, last_ping_ok, announced_at FROM peers WHERE last_ping_ok = TRUE AND http_url <> '' ORDER BY last_seen DESC NULLS LAST, did LIMIT $1", ) - .bind(limit.clamp(0, 201)) + .bind(limit.clamp(0, crate::api::repos::MAX_FEDERATED_PEERS as i64 + 1)) .fetch_all(&self.pool) .await?; Ok(rows @@ -8300,6 +8300,7 @@ mod peer_authority_tests { /// | `a_legacy_row_can_still_refresh_its_liveness` (db/mod.rs) | test-only. Seeds a PRE-GATE row by raw SQL on purpose: `upsert_peer` cannot create one, since the gate it is testing refuses exactly that DID. The fixture models what a deployed table already holds | /// | `gossip_ping_round_requires_two_failures_before_persisting_unreachable` (main.rs) | test-only fixture seed. Raw SQL because the test drives the readiness HYSTERESIS, which needs a row already at `last_ping_ok = TRUE` before the round runs; it never exercises the announce gate | /// | `manual_ping_uses_readiness_without_mutating_federation_gate` (api/peers.rs) | test-only fixture seed, same shape and same reason: the row under test must pre-exist so the assertion is about what the ping does NOT rewrite | +/// | `federated_peer_query_is_bounded` (api/repos.rs) | test-only fixture seed for the bounded federation query; inserts reachable rows without exercising peer admission | /// /// And the `upsert_peer` CALL-SITE authority table, which the ledger above /// structurally cannot hold, because the bootstrap site issues no SQL of its own @@ -8385,6 +8386,7 @@ mod peers_table_writer_guard { /// listed function that no longer has one. const LEDGER: &[(&str, usize)] = &[ ("a_legacy_row_can_still_refresh_its_liveness", 1), + ("federated_peer_query_is_bounded", 1), ( "gossip_ping_round_requires_two_failures_before_persisting_unreachable", 1, diff --git a/crates/gitlawb-node/src/main.rs b/crates/gitlawb-node/src/main.rs index 66bfa096..b3938bb5 100644 --- a/crates/gitlawb-node/src/main.rs +++ b/crates/gitlawb-node/src/main.rs @@ -396,6 +396,8 @@ async fn main() -> Result<()> { std::time::Duration::from_secs(3600), 200_000, ); + let federated_rate_limiter = + rate_limit::RateLimiter::new_bounded(12, std::time::Duration::from_secs(60), 10_000); if config.sync_trigger_rate_limit == 0 { tracing::warn!( "GITLAWB_SYNC_TRIGGER_RATE_LIMIT=0 — /sync/trigger IP rate limiting disabled" @@ -442,6 +444,7 @@ async fn main() -> Result<()> { push_limiter_trust, sync_trigger_rate_limiter, peer_write_rate_limiter, + federated_rate_limiter, shutdown_tx: shutdown_tx.clone(), git_read_semaphore: Arc::new(tokio::sync::Semaphore::new(config.max_concurrent_git_ops)), git_write_semaphore: Arc::new(tokio::sync::Semaphore::new( @@ -1189,6 +1192,7 @@ mod rate_limiter_sweep_tests { state.push_rate_limiter = RateLimiter::new(10, window); state.sync_trigger_rate_limiter = RateLimiter::new(10, window); state.peer_write_rate_limiter = RateLimiter::new(10, window); + state.federated_rate_limiter = RateLimiter::new(10, window); state.ipfs_rate_limiter = RateLimiter::new(10, window); state.ipfs_work_rate_limiter = RateLimiter::new(10, window); @@ -1199,6 +1203,7 @@ mod rate_limiter_sweep_tests { s.push_rate_limiter.clone(), s.sync_trigger_rate_limiter.clone(), s.peer_write_rate_limiter.clone(), + s.federated_rate_limiter.clone(), s.ipfs_rate_limiter.clone(), s.ipfs_work_rate_limiter.clone(), ] diff --git a/crates/gitlawb-node/src/server.rs b/crates/gitlawb-node/src/server.rs index cd38e55b..6b2ee45b 100644 --- a/crates/gitlawb-node/src/server.rs +++ b/crates/gitlawb-node/src/server.rs @@ -346,11 +346,7 @@ pub fn build_router(state: AppState) -> Router { get(repos::list_federated_repos) .route_layer(middleware::from_fn(rate_limit::rate_limit_by_ip)) .route_layer(axum::Extension(rate_limit::IpRateLimiter { - limiter: rate_limit::RateLimiter::new_bounded( - 12, - std::time::Duration::from_secs(60), - 10_000, - ), + limiter: state.federated_rate_limiter.clone(), trust: state.push_limiter_trust, })), ) diff --git a/crates/gitlawb-node/src/state.rs b/crates/gitlawb-node/src/state.rs index 24607e5a..062ad0dd 100644 --- a/crates/gitlawb-node/src/state.rs +++ b/crates/gitlawb-node/src/state.rs @@ -186,6 +186,9 @@ pub struct AppState { /// sink as trigger and accepts unsigned requests from known peers, so it is /// braked too; each peer's distinct IP gets its own bucket. pub peer_write_rate_limiter: RateLimiter, + /// Per-client-IP federation budget: 12 requests per minute, with at most + /// 10,000 tracked source keys. Shared by router instances and swept below. + pub federated_rate_limiter: RateLimiter, /// Process-wide graceful-shutdown signal. Sending `true` causes every /// task that holds a `watch::Receiver` to exit at its next await point. /// Used by: @@ -346,6 +349,7 @@ impl AppState { self.ipfs_work_rate_limiter.cleanup().await; self.sync_trigger_rate_limiter.cleanup().await; self.peer_write_rate_limiter.cleanup().await; + self.federated_rate_limiter.cleanup().await; } /// Trigger graceful shutdown. Idempotent — calling more than once diff --git a/crates/gitlawb-node/src/test_support.rs b/crates/gitlawb-node/src/test_support.rs index 430c0600..3edee88e 100644 --- a/crates/gitlawb-node/src/test_support.rs +++ b/crates/gitlawb-node/src/test_support.rs @@ -115,6 +115,7 @@ fn build_state(db: Arc, pool: PgPool) -> AppState { push_limiter_trust: crate::rate_limit::TrustedProxy::None, sync_trigger_rate_limiter: RateLimiter::new(60, Duration::from_secs(3600)), peer_write_rate_limiter: RateLimiter::new(600, Duration::from_secs(3600)), + federated_rate_limiter: RateLimiter::new_bounded(12, Duration::from_secs(60), 10_000), shutdown_tx: tokio::sync::watch::channel(false).0, // Generous — no test drives the handler-level shed (git_permit is unit-tested). git_read_semaphore: Arc::new(tokio::sync::Semaphore::new(64)),