Skip to content

Commit 6fa67d9

Browse files
committed
fix: address second re-review findings on bug_fix_2
- [P1] Add association-level row bound (MAX_ASSOC=2000) to SQL outer query, and 3-tuple cursor (pinned_at, sha256_hex, repo) for precise resumption when the bound is hit; all cursor advancement sites, opaque token, and sentinel updated to 3-tuple - [P1] Drain git stderr concurrently in run_git and assert_all_refs_are_commits to prevent deadlock when stderr fills - [P3] Move zero-limit fast path ahead of per-DID check so ?limit=0 requests never consume the per-DID quota - [P2] Short-circuit arweave anchor listing at limit=0 before repo enumeration, and add global rate-limit check - [P3] Include ipfs_list_rate_limiter and ipfs_list_global_limiter in the periodic cleanup task
1 parent 5d9d04f commit 6fa67d9

5 files changed

Lines changed: 126 additions & 68 deletions

File tree

‎crates/gitlawb-node/src/api/arweave.rs‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -87,6 +87,21 @@ pub async fn list_anchors(
8787
})));
8888
}
8989

90+
// Short-circuit for zero limit before any work or admission check (P2).
91+
if limit == 0 {
92+
return Ok(Json(serde_json::json!({
93+
"anchors": [],
94+
"count": 0,
95+
})));
96+
}
97+
98+
// Global rate limit: prevent DID-rotation enumeration (P2).
99+
if !state.ipfs_list_global_limiter.check("global").await {
100+
return Err(AppError::TooManyRequests(
101+
"rate limit exceeded for anchor listing".into(),
102+
));
103+
}
104+
90105
// Authenticated caller without ?repo=: scope to readable repos (P1).
91106
// Use the deduped, quarantine-filtered view (same as the pin listing).
92107
let repos = state

‎crates/gitlawb-node/src/api/ipfs.rs‎

Lines changed: 60 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -293,7 +293,7 @@ fn create_opaque_cursor(seed: &[u8; 32], cursor: &str) -> String {
293293

294294
/// Decode and verify an opaque truncated cursor token.
295295
/// Returns the original cursor string if valid and not expired.
296-
fn decode_opaque_cursor(seed: &[u8; 32], token: &str) -> Option<(String, String)> {
296+
fn decode_opaque_cursor(seed: &[u8; 32], token: &str) -> Option<(String, String, String)> {
297297
let cursor_key = derive_cursor_key(seed);
298298
let cipher = XChaCha20Poly1305::new_from_slice(&cursor_key)
299299
.expect("32-byte key is valid for XChaCha20Poly1305");
@@ -322,9 +322,10 @@ fn decode_opaque_cursor(seed: &[u8; 32], token: &str) -> Option<(String, String)
322322

323323
let cursor = std::str::from_utf8(&plaintext[8..]).ok()?;
324324

325-
let parts: Vec<&str> = cursor.splitn(2, '|').collect();
326-
if parts.len() == 2 {
327-
Some((parts[0].to_string(), parts[1].to_string()))
325+
let parts: Vec<&str> = cursor.splitn(3, '|').collect();
326+
if parts.len() >= 2 {
327+
let repo = parts.get(2).map(|s| s.to_string()).unwrap_or_default();
328+
Some((parts[0].to_string(), parts[1].to_string(), repo))
328329
} else {
329330
None
330331
}
@@ -437,15 +438,7 @@ pub async fn list_pins(
437438
let caller_str = caller.unwrap();
438439
let caller_owned = Some(caller_str.to_string());
439440

440-
// Per-DID rate limit: the listing performs expensive git walks and cat-file
441-
// probes, so a throwaway DID with a valid signature can exhaust resources (P1).
442-
if !state.ipfs_list_rate_limiter.check(caller_str).await {
443-
return Err(AppError::TooManyRequests(
444-
"rate limit exceeded for IPFS pin listing".into(),
445-
));
446-
}
447-
448-
// Clamp/handle zero limit before any expensive work so short-lived
441+
// Clamp/handle zero limit before any quota or expensive work so short-lived
449442
// requests do not drain rate-limit buckets or enumerate the node (P2).
450443
let max_visible = query.limit.clamp(0, 200);
451444

@@ -456,6 +449,14 @@ pub async fn list_pins(
456449
})));
457450
}
458451

452+
// Per-DID rate limit: the listing performs expensive git walks and cat-file
453+
// probes, so a throwaway DID with a valid signature can exhaust resources (P1).
454+
if !state.ipfs_list_rate_limiter.check(caller_str).await {
455+
return Err(AppError::TooManyRequests(
456+
"rate limit exceeded for IPFS pin listing".into(),
457+
));
458+
}
459+
459460
// Global rate limit: keyed on a fixed value so DID rotation cannot
460461
// bypass the enumeration cost guard (P1). Checked after the limit=0
461462
// fast path but before the whole-node reads so the bucket protects
@@ -497,8 +498,10 @@ pub async fn list_pins(
497498
}
498499

499500
// Decode the optional keyset cursor from base64.
500-
// Internal format: "pinned_at|sha256_hex" (2-tuple) for normal
501-
// pagination, or just "pinned_at" (1-tuple) for the truncated resume.
501+
// Wire format: "pinned_at|sha256_hex" (2-tuple). On resume we widen it
502+
// to (pinned_at, sha256_hex, "") so the SQL 3-tuple predicate skips
503+
// already-emitted SHAs while still processing remaining associations
504+
// of the batch's last SHA (P1).
502505
let decode_cursor = |s: &str| -> Option<(String, String)> {
503506
let bytes = URL_SAFE_NO_PAD.decode(s.as_bytes()).ok()?;
504507
let decoded = String::from_utf8(bytes).ok()?;
@@ -512,9 +515,9 @@ pub async fn list_pins(
512515
let encode_cursor =
513516
|pa: &str, sha: &str| -> String { URL_SAFE_NO_PAD.encode(format!("{pa}|{sha}")) };
514517

515-
let initial_cursor = match query.cursor.as_ref() {
518+
let initial_cursor: Option<(String, String, String)> = match query.cursor.as_ref() {
516519
Some(c) => match decode_cursor(c) {
517-
Some(cursor) => Some(cursor),
520+
Some((pa, sha)) => Some((pa, sha, String::new())),
518521
None => {
519522
return Err(AppError::BadRequest(
520523
"invalid cursor: expected base64-encoded pinned_at|sha256_hex".into(),
@@ -525,15 +528,15 @@ pub async fn list_pins(
525528
};
526529

527530
// Truncated resume cursor: XChaCha20Poly1305 AEAD token. Decrypts to the
528-
// same (pinned_at, sha256_hex) cursor on the server side but the
531+
// same (pinned_at, sha256_hex, repo) cursor on the server side but the
529532
// caller cannot decode hidden-row metadata from the wire format. If the
530533
// token is present but undecodable we return an explicit error so the
531534
// client does not silently restart at page 1.
532-
let truncated_resume = match query.truncated_cursor.as_ref() {
535+
let truncated_resume: Option<(String, String, String)> = match query.truncated_cursor.as_ref() {
533536
Some(t) => {
534537
let seed = state.node_keypair.to_seed();
535538
match decode_opaque_cursor(&seed, t) {
536-
Some(c) => Some(c),
539+
Some((pa, sha, repo)) => Some((pa, sha, repo)),
537540
None => {
538541
return Err(AppError::BadRequest(
539542
"invalid or expired truncated_cursor".into(),
@@ -567,7 +570,8 @@ pub async fn list_pins(
567570
// next_cursor is derived from the last *accepted* (visible) pin, never
568571
// from the last scanned row, to avoid leaking withheld-blob metadata
569572
// or skipping rows the caller was never shown.
570-
const BATCH_SIZE: i64 = 200;
573+
const BATCH_SIZE: i64 = 200; // max unique SHAs per batch
574+
const MAX_ASSOC: i64 = 2000; // max association rows per batch (SHA limit × 10)
571575
const MAX_BATCHES: usize = 10;
572576
const MAX_WALKS: usize = 50;
573577
const MAX_PROBES: usize = 200;
@@ -585,7 +589,7 @@ pub async fn list_pins(
585589
// repo associations. Track seen SHAs so each object is emitted at most
586590
// once, after evaluating per-association visibility (P2).
587591
let mut seen_shas: HashSet<String> = HashSet::new();
588-
let mut db_cursor: Option<(String, String)> = truncated_resume.or(initial_cursor);
592+
let mut db_cursor: Option<(String, String, String)> = truncated_resume.or(initial_cursor);
589593
let mut response_cursor: Option<(String, String)> = None;
590594
let mut allowed_blobs_by_repo: HashMap<String, (HashSet<String>, PathBuf)> = HashMap::new();
591595
let mut page_truncated = false;
@@ -610,9 +614,10 @@ pub async fn list_pins(
610614
&query_repos,
611615
&query_owner_dids,
612616
BATCH_SIZE,
617+
MAX_ASSOC,
613618
db_cursor
614619
.as_ref()
615-
.map(|(pa, sha)| (pa.as_str(), sha.as_str())),
620+
.map(|(pa, sha, repo)| (pa.as_str(), sha.as_str(), repo.as_str())),
616621
)
617622
.await
618623
.map_err(AppError::Internal)?
@@ -635,13 +640,17 @@ pub async fn list_pins(
635640

636641
for (i, pin) in batch.iter().enumerate() {
637642
if pin.repo.is_empty() {
638-
db_cursor = Some((pin.pinned_at.clone(), pin.sha256_hex.clone()));
643+
db_cursor = Some((pin.pinned_at.clone(), pin.sha256_hex.clone(), String::new()));
639644
pin_outcome.push(None);
640645
continue;
641646
}
642647
let Some((repo, rules)) = repos_by_slug.get(&pin.repo) else {
643648
// Unknown slug — advance cursor past it, no visibility check.
644-
db_cursor = Some((pin.pinned_at.clone(), pin.sha256_hex.clone()));
649+
db_cursor = Some((
650+
pin.pinned_at.clone(),
651+
pin.sha256_hex.clone(),
652+
pin.repo.clone(),
653+
));
645654
pin_outcome.push(None);
646655
continue;
647656
};
@@ -777,7 +786,11 @@ pub async fn list_pins(
777786
// All pins had empty/unmatched repos — advance past the batch so
778787
// we don't loop forever on the same unprocessable rows (P1).
779788
if let Some(last) = batch.last() {
780-
db_cursor = Some((last.pinned_at.clone(), last.sha256_hex.clone()));
789+
db_cursor = Some((
790+
last.pinned_at.clone(),
791+
last.sha256_hex.clone(),
792+
last.repo.clone(),
793+
));
781794
}
782795
}
783796

@@ -895,37 +908,47 @@ pub async fn list_pins(
895908
if batch_cursor.is_some() {
896909
db_cursor = batch_cursor;
897910
}
898-
} else if let Some(pin) = i.checked_sub(1).and_then(|prev| batch.get(prev)) {
899-
db_cursor = Some((pin.pinned_at.clone(), pin.sha256_hex.clone()));
911+
} else if let Some(prev) = i.checked_sub(1).and_then(|prev| batch.get(prev)) {
912+
db_cursor = Some((
913+
prev.pinned_at.clone(),
914+
prev.sha256_hex.clone(),
915+
prev.repo.clone(),
916+
));
900917
}
901918
break;
902919
}
903920

904921
let pin = batch[i].clone();
905922
let Some((repo, rules)) = repos_by_slug.get(&pin.repo) else {
906923
// Already advanced past in phase 1 — just maintain cursor.
907-
db_cursor = Some((pin.pinned_at.clone(), pin.sha256_hex.clone()));
924+
db_cursor = Some((
925+
pin.pinned_at.clone(),
926+
pin.sha256_hex.clone(),
927+
pin.repo.clone(),
928+
));
908929
continue;
909930
};
910931

911932
if !has_path_scoped_rule(rules) {
912933
let pa = pin.pinned_at.clone();
913934
let sha = pin.sha256_hex.clone();
935+
let repo_slug = pin.repo.clone();
914936

915937
// Dedup by sha256_hex (P2): only skip if already *emitted* —
916938
// do not suppress a visible association because a hidden one
917939
// appeared first in the batch.
918940
if !seen_shas.insert(sha.clone()) {
919-
db_cursor = Some((pa, sha));
941+
db_cursor = Some((pa, sha, repo_slug));
920942
continue;
921943
}
922944

923945
response_cursor = Some((pa.clone(), sha.clone()));
924946
pins.push(pin);
925-
db_cursor = Some((pa, sha));
947+
db_cursor = Some((pa, sha, repo_slug));
926948
} else {
927949
let pa = pin.pinned_at.clone();
928950
let sha = pin.sha256_hex.clone();
951+
let repo_slug = pin.repo.clone();
929952

930953
let visible = match pin_outcome[i] {
931954
Some(v) => v,
@@ -943,13 +966,13 @@ pub async fn list_pins(
943966
};
944967
if visible {
945968
if !seen_shas.insert(sha.clone()) {
946-
db_cursor = Some((pa, sha));
969+
db_cursor = Some((pa, sha, repo_slug));
947970
continue;
948971
}
949972
response_cursor = Some((pa.clone(), sha.clone()));
950973
pins.push(pin);
951974
}
952-
db_cursor = Some((pa, sha));
975+
db_cursor = Some((pa, sha, repo_slug));
953976
}
954977

955978
if pins.len() >= max_visible as usize {
@@ -966,10 +989,10 @@ pub async fn list_pins(
966989
// When page 1 is all-deferred (no walk permit available) neither
967990
// response_cursor nor db_cursor was set. Emit a sentinel opaque cursor
968991
// so the client can retry; on the retry the sentinel decodes to
969-
// ("\x7f", "\x7f") which the keyset WHERE < predicate treats
992+
// ("\x7f", "\x7f", "\x7f") which the keyset WHERE < predicate treats
970993
// as "include every row" — effectively restarting from the beginning (P2).
971994
if page_truncated && response_cursor.is_none() && db_cursor.is_none() {
972-
db_cursor = Some(("\x7f".to_string(), "\x7f".to_string()));
995+
db_cursor = Some(("\x7f".to_string(), "\x7f".to_string(), "\x7f".to_string()));
973996
}
974997

975998
let mut body = serde_json::json!({
@@ -991,8 +1014,8 @@ pub async fn list_pins(
9911014
// when there are no visible rows to derive a keyset cursor from.
9921015
if let Some((ref pa, ref sha)) = response_cursor {
9931016
body["next_cursor"] = serde_json::json!(encode_cursor(pa, sha));
994-
} else if let Some((ref pa, ref sha)) = db_cursor {
995-
let cursor_str = format!("{pa}|{sha}");
1017+
} else if let Some((ref pa, ref sha, ref repo)) = db_cursor {
1018+
let cursor_str = format!("{pa}|{sha}|{repo}");
9961019
let seed = state.node_keypair.to_seed();
9971020
let token = create_opaque_cursor(&seed, &cursor_str);
9981021
body["truncated_cursor"] = serde_json::json!(token);

‎crates/gitlawb-node/src/db/mod.rs‎

Lines changed: 24 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -2473,21 +2473,23 @@ impl Db {
24732473
}
24742474

24752475
/// Bounded global pin query: returns pins for any of the given (repo, owner_did)
2476-
/// pairs, ordered by pinned_at DESC, capped at `limit` distinct objects.
2477-
/// Each batch returns up to `limit` unique SHAs (determined by DISTINCT ON)
2478-
/// and ALL repo associations for those SHAs, so the caller can evaluate
2479-
/// per-association visibility before deduplicating (P2).
2480-
/// Uses keyset pagination on (pinned_at, sha256_hex). Resume from a client
2481-
/// 2-tuple cursor uses the same (pinned_at, sha256_hex) directly; all
2482-
/// associations of the last SHA are guaranteed to be in the same batch.
2476+
/// pairs, ordered by (pinned_at DESC, sha256_hex DESC, repo DESC) and capped
2477+
/// by both distinct-object count and total-association row count.
2478+
/// The inner DISTINCT ON picks up to `sha_limit` unique SHAs; the outer query
2479+
/// then returns ALL repo associations for those SHAs, bounded by `assoc_limit`
2480+
/// rows. When `cursor` is `Some((pinned_at, sha256_hex, repo))`, the inner
2481+
/// query skips rows whose keyset position is not strictly before the cursor,
2482+
/// so a partially-returned SHA (outer result hit `assoc_limit`) resumes
2483+
/// correctly on the next fetch.
24832484
pub async fn list_pinned_cids_for_repos(
24842485
&self,
24852486
repos: &[String],
24862487
owner_dids: &[String],
2487-
limit: i64,
2488-
cursor: Option<(&str, &str)>,
2488+
sha_limit: i64,
2489+
assoc_limit: i64,
2490+
cursor: Option<(&str, &str, &str)>,
24892491
) -> Result<Vec<PinnedCidRecord>> {
2490-
let rows = if let Some((pa, sha)) = cursor {
2492+
let rows = if let Some((pa, sha, repo)) = cursor {
24912493
sqlx::query(
24922494
r#"WITH batch_shas AS (
24932495
SELECT sha256_hex
@@ -2499,11 +2501,12 @@ impl Db {
24992501
AS pairs(repo, owner_did)
25002502
ON (COALESCE(pr.repo, p.repo), COALESCE(pr.owner_did, p.owner_did))
25012503
= (pairs.repo, pairs.owner_did)
2502-
WHERE (p.pinned_at, p.sha256_hex) < ($3::text, $4::text)
2504+
WHERE (p.pinned_at, p.sha256_hex, COALESCE(pr.repo, p.repo))
2505+
< ($3::text, $4::text, $5::text)
25032506
ORDER BY p.sha256_hex, p.pinned_at DESC
25042507
) deduped
25052508
ORDER BY pinned_at DESC, sha256_hex DESC
2506-
LIMIT $5
2509+
LIMIT $6
25072510
)
25082511
SELECT p.sha256_hex, p.cid, p.pinned_at, p.pinata_cid,
25092512
COALESCE(pr.repo, p.repo) AS repo,
@@ -2516,13 +2519,16 @@ impl Db {
25162519
ON (COALESCE(pr.repo, p.repo), COALESCE(pr.owner_did, p.owner_did))
25172520
= (pairs.repo, pairs.owner_did)
25182521
ORDER BY p.pinned_at DESC, p.sha256_hex DESC,
2519-
COALESCE(pr.repo, p.repo) DESC"#,
2522+
COALESCE(pr.repo, p.repo) DESC
2523+
LIMIT $7"#,
25202524
)
25212525
.bind(repos)
25222526
.bind(owner_dids)
25232527
.bind(pa)
25242528
.bind(sha)
2525-
.bind(limit)
2529+
.bind(repo)
2530+
.bind(sha_limit)
2531+
.bind(assoc_limit)
25262532
.fetch_all(&self.pool)
25272533
.await?
25282534
} else {
@@ -2553,11 +2559,13 @@ impl Db {
25532559
ON (COALESCE(pr.repo, p.repo), COALESCE(pr.owner_did, p.owner_did))
25542560
= (pairs.repo, pairs.owner_did)
25552561
ORDER BY p.pinned_at DESC, p.sha256_hex DESC,
2556-
COALESCE(pr.repo, p.repo) DESC"#,
2562+
COALESCE(pr.repo, p.repo) DESC
2563+
LIMIT $4"#,
25572564
)
25582565
.bind(repos)
25592566
.bind(owner_dids)
2560-
.bind(limit)
2567+
.bind(sha_limit)
2568+
.bind(assoc_limit)
25612569
.fetch_all(&self.pool)
25622570
.await?
25632571
};

0 commit comments

Comments
 (0)