A thin push reads at most 200 delta bases from the store, 16 at a time, and stops when the store says it is busy, so a large push no longer fans out into hundreds of reads at once that trip the store's breaker and come back 503; the checks say how many bases the pack lacked and how many were read.
2 files+214−330/2 viewed
Binary or large file; its contents are not shown.
| 41 | 41 | const MAX_PUSH_COMMITS: usize = 300; | |
| 42 | 42 | /// Files compared per commit, at most. | |
| 43 | 43 | const MAX_FILES_PER_COMMIT: usize = 300; | |
| 44 | − | /// Bases fetched from the store for a thin pack, at most. | |
| 45 | − | const MAX_BASES: usize = 500; | |
| 44 | + | /// Bases fetched from the store for a thin pack, at most, in all rounds | |
| 45 | + | /// together. Each is one or two store reads, each with a Cache API look. | |
| 46 | + | const MAX_BASES: usize = 200; | |
| 47 | + | /// Bases asked for at a time (two reads each when the pack does not say | |
| 48 | + | /// whether a base is a blob or a tree). | |
| 49 | + | const BASES_AT_ONCE: usize = 16; | |
| 46 | 50 | /// The largest push that is read whole and scanned. A larger one is | |
| 47 | 51 | /// declined, since it cannot be checked (git_http.rs `LargePushes`). | |
| 48 | 52 | pub const MAX_SCANNED_PUSH: usize = 24 * 1024 * 1024; | |
| ⋯ | |||
| 260 | 264 | Ok(found) | |
| 261 | 265 | } | |
| 262 | 266 | ||
| 263 | − | /// Fetches what a thin pack's deltas are based on from the repository, | |
| 264 | − | /// all at once. A base the pack's own trees name is read as what they say | |
| 265 | − | /// it is; any other is asked for as a blob and as a tree together, and | |
| 266 | − | /// whichever it is answers. | |
| 267 | − | pub(crate) async fn supply_bases<R: GitRepo>(pack: &mut Pack, repo: &R) -> Result<()> { | |
| 267 | + | /// What [`supply_bases`] did for a pack. | |
| 268 | + | #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] | |
| 269 | + | pub struct Bases { | |
| 270 | + | /// Bases the pack's deltas needed from outside it when it arrived: 0 | |
| 271 | + | /// for a pack sent whole, as g1t asks for (`no-thin`, git_http.rs). | |
| 272 | + | pub missing: usize, | |
| 273 | + | /// Bases asked of the store, in every round. | |
| 274 | + | pub asked: usize, | |
| 275 | + | /// Bases still missing at the end: objects that stay unresolved. | |
| 276 | + | pub left: usize, | |
| 277 | + | } | |
| 278 | + | ||
| 279 | + | impl Bases { | |
| 280 | + | /// Whether the pack came thin, its deltas based on objects outside it. | |
| 281 | + | pub fn thin(&self) -> bool { | |
| 282 | + | self.missing > 0 | |
| 283 | + | } | |
| 284 | + | } | |
| 285 | + | ||
| 286 | + | /// Fetches what a thin pack's deltas are based on from the repository. | |
| 287 | + | /// A base the pack's own trees name is read as what they say it is; any | |
| 288 | + | /// other is asked for as a blob and as a tree together, and whichever it | |
| 289 | + | /// is answers. | |
| 290 | + | /// | |
| 291 | + | /// g1t asks for packs without outside bases (`no-thin`), so this is for | |
| 292 | + | /// clients that send thin ones anyway, and it is bounded: [`MAX_BASES`] in | |
| 293 | + | /// all, over up to three rounds for delta chains, [`BASES_AT_ONCE`] at a | |
| 294 | + | /// time. Asking for hundreds at once made a large thin push fan out into | |
| 295 | + | /// as many store reads, and Cache API looks beside them, all in flight | |
| 296 | + | /// together; the reads that failed were tried again and counted against | |
| 297 | + | /// the store's breaker (store.rs `invoke`), which then turned every read | |
| 298 | + | /// after them away, so the push was answered 503. A store that says it is | |
| 299 | + | /// busy stops the reading here. Bases left missing leave their objects | |
| 300 | + | /// unresolved, and the checks go on without them, as for any pack the | |
| 301 | + | /// store cannot complete. | |
| 302 | + | pub(crate) async fn supply_bases<R: GitRepo>(pack: &mut Pack, repo: &R) -> Result<Bases> { | |
| 303 | + | let mut bases = Bases { missing: pack.missing_bases().len(), ..Bases::default() }; | |
| 304 | + | let mut busy = false; | |
| 268 | 305 | for _ in 0..3 { | |
| 269 | 306 | let missing = pack.missing_bases(); | |
| 270 | − | if missing.is_empty() { | |
| 271 | − | return Ok(()); | |
| 307 | + | if missing.is_empty() || busy || bases.asked >= MAX_BASES { | |
| 308 | + | break; | |
| 272 | 309 | } | |
| 273 | 310 | let named = pack.named_kinds(); | |
| 274 | − | let blob = async |id: &str| repo.read_blob(id).await.ok().flatten().map(|bytes| (ObjectKind::Blob, bytes)); | |
| 311 | + | // A read that fails is a base not found, unless the store is busy, | |
| 312 | + | // which ends the reading. | |
| 313 | + | let settle = |read: Result<Option<(ObjectKind, Vec<u8>)>>| match read { | |
| 314 | + | Ok(found) => Ok(found), | |
| 315 | + | Err(error) if crate::resilience::busy(&error.to_string()).is_some() => Err(()), | |
| 316 | + | Err(_) => Ok(None), | |
| 317 | + | }; | |
| 318 | + | let blob = async |id: &str| settle(repo.read_blob(id).await.map(|found| found.map(|bytes| (ObjectKind::Blob, bytes)))); | |
| 275 | 319 | let tree = async |id: &str| { | |
| 276 | − | repo.read_tree(id).await.ok().flatten().map(|entries| { | |
| 277 | − | let items: Vec<TreeItem> = entries | |
| 278 | − | .into_iter() | |
| 279 | − | .map(|entry| TreeItem { mode: mode(entry.kind).to_owned(), name: entry.name, id: entry.hash }) | |
| 280 | − | .collect(); | |
| 281 | − | (ObjectKind::Tree, encode_tree(&items)) | |
| 282 | − | }) | |
| 320 | + | settle(repo.read_tree(id).await.map(|found| { | |
| 321 | + | found.map(|entries| { | |
| 322 | + | let items: Vec<TreeItem> = entries | |
| 323 | + | .into_iter() | |
| 324 | + | .map(|entry| TreeItem { mode: mode(entry.kind).to_owned(), name: entry.name, id: entry.hash }) | |
| 325 | + | .collect(); | |
| 326 | + | (ObjectKind::Tree, encode_tree(&items)) | |
| 327 | + | }) | |
| 328 | + | })) | |
| 283 | 329 | }; | |
| 284 | − | let found = futures_util::future::join_all(missing.iter().take(MAX_BASES).map(|id| async { | |
| 285 | − | match named.get(id.as_str()) { | |
| 286 | − | Some(ObjectKind::Tree) => tree(id).await, | |
| 287 | − | Some(_) => blob(id).await, | |
| 288 | − | None => { | |
| 289 | − | let (as_blob, as_tree) = futures_util::future::join(blob(id), tree(id)).await; | |
| 290 | − | as_blob.or(as_tree) | |
| 330 | + | let wanted: Vec<&String> = missing.iter().take(MAX_BASES - bases.asked).collect(); | |
| 331 | + | let mut progress = false; | |
| 332 | + | for batch in wanted.chunks(BASES_AT_ONCE) { | |
| 333 | + | let found = futures_util::future::join_all(batch.iter().map(|id| async { | |
| 334 | + | match named.get(id.as_str()) { | |
| 335 | + | Some(ObjectKind::Tree) => tree(id).await, | |
| 336 | + | Some(_) => blob(id).await, | |
| 337 | + | None => { | |
| 338 | + | let (as_blob, as_tree) = futures_util::future::join(blob(id), tree(id)).await; | |
| 339 | + | match (as_blob, as_tree) { | |
| 340 | + | (Ok(Some(found)), _) | (_, Ok(Some(found))) => Ok(Some(found)), | |
| 341 | + | (Err(()), _) | (_, Err(())) => Err(()), | |
| 342 | + | _ => Ok(None), | |
| 343 | + | } | |
| 344 | + | } | |
| 291 | 345 | } | |
| 346 | + | })) | |
| 347 | + | .await; | |
| 348 | + | bases.asked += batch.len(); | |
| 349 | + | for (id, object) in batch.iter().zip(found) { | |
| 350 | + | match object { | |
| 351 | + | Ok(Some((kind, data))) => { | |
| 352 | + | pack.supply(id, kind, data); | |
| 353 | + | progress = true; | |
| 354 | + | } | |
| 355 | + | Ok(None) => {} | |
| 356 | + | Err(()) => busy = true, | |
| 357 | + | } | |
| 292 | 358 | } | |
| 293 | − | })) | |
| 294 | − | .await; | |
| 295 | − | let mut progress = false; | |
| 296 | − | for (id, object) in missing.iter().zip(found) { | |
| 297 | − | if let Some((kind, data)) = object { | |
| 298 | − | pack.supply(id, kind, data); | |
| 299 | − | progress = true; | |
| 359 | + | if busy { | |
| 360 | + | break; | |
| 300 | 361 | } | |
| 301 | 362 | } | |
| 302 | 363 | if !progress { | |
| 303 | − | return Ok(()); | |
| 364 | + | break; | |
| 304 | 365 | } | |
| 305 | 366 | } | |
| 306 | − | Ok(()) | |
| 367 | + | bases.left = pack.missing_bases().len(); | |
| 368 | + | Ok(bases) | |
| 307 | 369 | } | |
| 308 | 370 | ||
| 309 | 371 | /// The secrets the commits in a push add, each secret once. A push too | |
| ⋯ | |||
| 760 | 822 | blob_reads: Cell<u32>, | |
| 761 | 823 | tree_reads: Cell<u32>, | |
| 762 | 824 | log_reads: Cell<u32>, | |
| 825 | + | /// Each read waits once before it answers, as a store's does, so | |
| 826 | + | /// that reads asked for together are in flight together. | |
| 827 | + | waits: bool, | |
| 828 | + | in_flight: Cell<u32>, | |
| 829 | + | most_in_flight: Cell<u32>, | |
| 830 | + | /// Every read fails as a busy store's does. | |
| 831 | + | busy: bool, | |
| 832 | + | } | |
| 833 | + | ||
| 834 | + | impl FakeRepo { | |
| 835 | + | async fn reading(&self) -> Result<()> { | |
| 836 | + | if self.busy { | |
| 837 | + | return Err(crate::resilience::Busy { rate_limited: true, retry_after: 5, read_only: false }.error("readBlob")); | |
| 838 | + | } | |
| 839 | + | if self.waits { | |
| 840 | + | self.in_flight.set(self.in_flight.get() + 1); | |
| 841 | + | self.most_in_flight.set(self.most_in_flight.get().max(self.in_flight.get())); | |
| 842 | + | YieldOnce(false).await; | |
| 843 | + | self.in_flight.set(self.in_flight.get() - 1); | |
| 844 | + | } | |
| 845 | + | Ok(()) | |
| 846 | + | } | |
| 847 | + | } | |
| 848 | + | ||
| 849 | + | /// Waits once, then is ready. | |
| 850 | + | struct YieldOnce(bool); | |
| 851 | + | ||
| 852 | + | impl Future for YieldOnce { | |
| 853 | + | type Output = (); | |
| 854 | + | fn poll(mut self: std::pin::Pin<&mut Self>, context: &mut Context<'_>) -> Poll<()> { | |
| 855 | + | if self.0 { | |
| 856 | + | return Poll::Ready(()); | |
| 857 | + | } | |
| 858 | + | self.0 = true; | |
| 859 | + | context.waker().wake_by_ref(); | |
| 860 | + | Poll::Pending | |
| 861 | + | } | |
| 763 | 862 | } | |
| 764 | 863 | ||
| 864 | + | /// Runs a future to its end, polling it until it is ready. | |
| 865 | + | fn run_waiting<F: Future>(future: F) -> F::Output { | |
| 866 | + | let mut future = pin!(future); | |
| 867 | + | loop { | |
| 868 | + | if let Poll::Ready(output) = future.as_mut().poll(&mut Context::from_waker(Waker::noop())) { | |
| 869 | + | return output; | |
| 870 | + | } | |
| 871 | + | } | |
| 872 | + | } | |
| 873 | + | ||
| 765 | 874 | impl GitRepo for FakeRepo { | |
| 766 | 875 | async fn access(&self, _scope: Scope) -> Result<GitAccess> { | |
| 767 | 876 | unimplemented!() | |
| ⋯ | |||
| 778 | 887 | } | |
| 779 | 888 | async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>> { | |
| 780 | 889 | self.tree_reads.set(self.tree_reads.get() + 1); | |
| 890 | + | self.reading().await?; | |
| 781 | 891 | Ok(self.trees.get(tree_hash).cloned()) | |
| 782 | 892 | } | |
| 783 | 893 | async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>> { | |
| 784 | 894 | self.blob_reads.set(self.blob_reads.get() + 1); | |
| 895 | + | self.reading().await?; | |
| 785 | 896 | Ok(self.blobs.get(blob_hash).cloned()) | |
| 786 | 897 | } | |
| 787 | 898 | async fn read_file(&self, _git_ref: &str, _path: &str) -> Result<Option<Vec<u8>>> { | |
| ⋯ | |||
| 968 | 1079 | } | |
| 969 | 1080 | ||
| 970 | 1081 | #[test] | |
| 1082 | + | fn a_pack_sent_whole_is_checked_without_reading_a_base() { | |
| 1083 | + | // What git sends when told `no-thin`: a delta's base is in the pack | |
| 1084 | + | // with it. Here the second file is a delta on the first. | |
| 1085 | + | let first = b"REGION=eu | |
| 1086 | + | ".to_vec(); | |
| 1087 | + | let first_id = object_id(ObjectKind::Blob, &first); | |
| 1088 | + | let added = format!("AWS_KEY={} | |
| 1089 | + | ", key()).into_bytes(); | |
| 1090 | + | let second = [first.clone(), added.clone()].concat(); | |
| 1091 | + | let mut delta = vec![first.len() as u8, second.len() as u8, 0x80 | 0x10, first.len() as u8, added.len() as u8]; | |
| 1092 | + | delta.extend_from_slice(&added); | |
| 1093 | + | let tree = encode_tree(&[ | |
| 1094 | + | TreeItem { mode: "100644".into(), name: "a.env".into(), id: first_id.clone() }, | |
| 1095 | + | TreeItem { mode: "100644".into(), name: "b.env".into(), id: object_id(ObjectKind::Blob, &second) }, | |
| 1096 | + | ]); | |
| 1097 | + | let tree_id = object_id(ObjectKind::Tree, &tree); | |
| 1098 | + | let body = push(&[ | |
| 1099 | + | Entry::Whole(ObjectKind::Commit, commit(&tree_id, None)), | |
| 1100 | + | Entry::Whole(ObjectKind::Tree, tree), | |
| 1101 | + | Entry::Whole(ObjectKind::Blob, first), | |
| 1102 | + | Entry::Delta(first_id, delta), | |
| 1103 | + | ]); | |
| 1104 | + | let repo = FakeRepo::default(); | |
| 1105 | + | let mut pack = crate::push_checks::read_pack(&body).unwrap(); | |
| 1106 | + | let bases = run(supply_bases(&mut pack, &repo)).unwrap(); | |
| 1107 | + | assert_eq!(bases, Bases::default()); | |
| 1108 | + | assert!(!bases.thin()); | |
| 1109 | + | assert_eq!(pack.unresolved(), 0); | |
| 1110 | + | let found = run(scan_pack(&repo, body.len(), &Ok(pack), &[])).unwrap(); | |
| 1111 | + | assert_eq!(found.len(), 1); | |
| 1112 | + | assert_eq!((found[0].path.as_str(), found[0].line), ("b.env", 2)); | |
| 1113 | + | // Nothing was read from the store: not a base, not a file. | |
| 1114 | + | assert_eq!((repo.blob_reads.get(), repo.tree_reads.get(), repo.log_reads.get()), (0, 0, 0)); | |
| 1115 | + | } | |
| 1116 | + | ||
| 1117 | + | /// A pack of `count` deltas, each on a different base outside it. | |
| 1118 | + | fn thin_pack(count: usize) -> Pack { | |
| 1119 | + | let entries: Vec<Entry> = (1..=count) | |
| 1120 | + | .map(|at| Entry::Delta(format!("{at:040x}"), vec![1, 2, 0x02, b'h', b'i'])) | |
| 1121 | + | .collect(); | |
| 1122 | + | crate::push_checks::read_pack(&push(&entries)).unwrap() | |
| 1123 | + | } | |
| 1124 | + | ||
| 1125 | + | #[test] | |
| 1126 | + | fn a_large_thin_push_reads_a_bounded_number_of_bases_a_few_at_a_time() { | |
| 1127 | + | // 300 bases the store does not have, none named by a tree: before, | |
| 1128 | + | // each round asked for 500 at once, two reads each, three rounds. | |
| 1129 | + | let mut pack = thin_pack(300); | |
| 1130 | + | let repo = FakeRepo { waits: true, ..FakeRepo::default() }; | |
| 1131 | + | let bases = run_waiting(supply_bases(&mut pack, &repo)).unwrap(); | |
| 1132 | + | assert_eq!(bases, Bases { missing: 300, asked: MAX_BASES, left: 300 }); | |
| 1133 | + | assert!(bases.thin()); | |
| 1134 | + | // Asked once each, as a blob and as a tree, and never more at once | |
| 1135 | + | // than a batch's. | |
| 1136 | + | assert_eq!((repo.blob_reads.get(), repo.tree_reads.get()), (MAX_BASES as u32, MAX_BASES as u32)); | |
| 1137 | + | assert_eq!(repo.most_in_flight.get(), 2 * BASES_AT_ONCE as u32); | |
| 1138 | + | } | |
| 1139 | + | ||
| 1140 | + | #[test] | |
| 1141 | + | fn a_busy_store_stops_the_reading_of_bases() { | |
| 1142 | + | let mut pack = thin_pack(100); | |
| 1143 | + | let repo = FakeRepo { busy: true, ..FakeRepo::default() }; | |
| 1144 | + | let bases = run(supply_bases(&mut pack, &repo)).unwrap(); | |
| 1145 | + | // One batch, then no more: the checks go on, and the store's own | |
| 1146 | + | // answer to them says it is busy. | |
| 1147 | + | assert_eq!(bases, Bases { missing: 100, asked: BASES_AT_ONCE, left: 100 }); | |
| 1148 | + | assert_eq!(repo.blob_reads.get(), BASES_AT_ONCE as u32); | |
| 1149 | + | } | |
| 1150 | + | ||
| 1151 | + | #[test] | |
| 971 | 1152 | fn the_commit_a_push_builds_on_is_read_once_however_many_commits_build_on_it() { | |
| 972 | 1153 | // Two branches pushed at once, each one commit on the same parent. | |
| 973 | 1154 | let parent_id = "c71546fcd893ef8b0f57388b65e620d759705dda".to_owned(); | |