1use std::collections::{BTreeMap, BTreeSet};
6
7use mkit_core::hash::Hash;
8use mkit_core::store::MAX_RAW_OBJECT_SIZE;
9
10use super::codec::{self, RelayV1};
11use super::keys::{self, ParsedKey};
12use super::outbox::MAX_RELAY_PUTS;
13use super::{
14 BlobKey, Key, MAX_BATCH_BYTES, MAX_KEY_BYTES, MAX_VALUE_BYTES, NamespaceStore, Partition,
15 RangeScan, StoreError, Value,
16};
17use crate::pipeline::ShardMap;
18use crate::repo::RepoId;
19
20pub const MAX_LOOKUP_IDS: usize = 256;
22pub const MAX_LOOKUP_ROWS: usize = 4096;
26pub const MAX_LOOKUP_PAGES: usize = 512;
29pub const MAX_LOOKUP_MEMBERSHIP_READS: usize = 487;
32const CANDIDATE_BYTES: usize = 4 * 1024 * 1024;
38const CANDIDATE_OVERHEAD: usize = 512;
39const MEMBERSHIP_CHUNK: usize = 128;
40const SCAN_PAGE_ROWS: u32 = 128;
41const SCAN_CALL_ROWS: u32 = 1_000;
42const _: () = assert!(MAX_LOOKUP_PAGES + MAX_LOOKUP_MEMBERSHIP_READS < 1000);
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub struct IndexValue {
47 pub frame_offset: u64,
49 pub frame_length: u64,
51 pub wire_type: u8,
53 pub decoded_size: u64,
55 pub chain_depth: u32,
59 pub delta_base: Option<Hash>,
61}
62
63impl IndexValue {
64 pub(crate) fn validate(&self, object: &Hash) -> Result<(), StoreError> {
65 if self.frame_length == 0
66 || self.frame_length > u64::from(u32::MAX) + 5
67 || self.frame_offset.checked_add(self.frame_length).is_none()
68 || self.decoded_size == 0
69 || self.decoded_size > MAX_RAW_OBJECT_SIZE as u64
70 || self.chain_depth > u32::from(u16::MAX)
71 || self.delta_base == Some(*object)
72 {
73 return Err(StoreError::Invalid("invalid object index metadata".into()));
74 }
75 match (self.wire_type, self.delta_base, self.chain_depth) {
76 (0x00 | 0x03, None, 0) | (0x02 | 0x04, Some(_), 1..) => Ok(()),
77 _ => Err(StoreError::Invalid(
78 "invalid object index entry type".into(),
79 )),
80 }
81 }
82}
83
84#[derive(Debug, Clone, Copy, PartialEq, Eq)]
87pub struct IndexEntry {
88 pub object: Hash,
90 pub value: IndexValue,
92}
93
94#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct DirectIndexBatch {
97 pub target: Partition,
99 pub puts: Vec<(Key, Value)>,
101}
102
103#[derive(Debug, Clone, Default, PartialEq, Eq)]
106pub struct IndexPlan {
107 pub direct: Vec<DirectIndexBatch>,
109 pub relay: Vec<RelayV1>,
111}
112
113pub fn plan_index_rows(
117 shards: &dyn ShardMap,
118 repo: &RepoId,
119 source: &Partition,
120 pack: &Hash,
121 entries: &[IndexEntry],
122 at_ms: u64,
123) -> Result<IndexPlan, StoreError> {
124 plan_index_rows_inner(shards, repo, source, pack, entries, at_ms, false)
125}
126
127pub fn plan_index_rows_direct(
131 shards: &dyn ShardMap,
132 repo: &RepoId,
133 source: &Partition,
134 pack: &Hash,
135 entries: &[IndexEntry],
136 at_ms: u64,
137) -> Result<IndexPlan, StoreError> {
138 plan_index_rows_inner(shards, repo, source, pack, entries, at_ms, true)
139}
140
141fn plan_index_rows_inner(
142 shards: &dyn ShardMap,
143 repo: &RepoId,
144 source: &Partition,
145 pack: &Hash,
146 entries: &[IndexEntry],
147 at_ms: u64,
148 direct_all: bool,
149) -> Result<IndexPlan, StoreError> {
150 let mut grouped: BTreeMap<Partition, BTreeMap<Key, Value>> = BTreeMap::new();
151 for entry in entries {
152 let target = shards.object_index(repo, &entry.object);
153 let key = keys::object_index(&repo.name, &entry.object, pack);
154 let group = grouped.entry(target).or_default();
155 if let std::collections::btree_map::Entry::Vacant(slot) = group.entry(key) {
156 slot.insert(codec::encode_object_index(&entry.object, &entry.value)?);
157 }
158 }
159 let mut plan = IndexPlan::default();
160 for (target, rows) in grouped {
161 if direct_all || &target == source {
162 let mut puts = Vec::new();
163 let mut bytes = 0;
164 for (key, value) in rows {
165 let size = key.as_bytes().len() + value.as_bytes().len();
166 if size > MAX_BATCH_BYTES
167 || key.as_bytes().len() > MAX_KEY_BYTES
168 || value.as_bytes().len() > MAX_VALUE_BYTES
169 {
170 return Err(StoreError::Invalid("object index row too large".into()));
171 }
172 if puts.len() == MAX_RELAY_PUTS || bytes + size > MAX_BATCH_BYTES {
173 plan.direct.push(DirectIndexBatch {
174 target: target.clone(),
175 puts: std::mem::take(&mut puts),
176 });
177 bytes = 0;
178 }
179 bytes += size;
180 puts.push((key, value));
181 }
182 if !puts.is_empty() {
183 plan.direct.push(DirectIndexBatch { target, puts });
184 }
185 } else {
186 let mut row = RelayV1 {
187 at_ms,
188 target: target.clone(),
189 puts: Vec::new(),
190 deletes: Vec::new(),
191 };
192 let base_bytes = codec::encode_relay(&row)?.as_bytes().len();
193 let mut encoded_bytes = base_bytes;
194 for (key, value) in rows {
195 let addition = 7 + 2 * (key.as_bytes().len() + value.as_bytes().len());
197 if key.as_bytes().len() > MAX_KEY_BYTES
198 || value.as_bytes().len() > MAX_VALUE_BYTES
199 || base_bytes + addition > MAX_VALUE_BYTES
200 || base_bytes + addition + 2 * (MAX_KEY_BYTES + 8) > MAX_BATCH_BYTES
201 {
202 return Err(StoreError::Invalid("object index row too large".into()));
203 }
204 let comma = usize::from(!row.puts.is_empty());
205 if row.puts.len() == MAX_RELAY_PUTS
206 || encoded_bytes + addition + comma > MAX_VALUE_BYTES
207 || encoded_bytes + addition + comma + 2 * (MAX_KEY_BYTES + 8) > MAX_BATCH_BYTES
208 {
209 codec::encode_relay(&row)?;
211 plan.relay.push(row);
212 row = RelayV1 {
213 at_ms,
214 target: target.clone(),
215 puts: Vec::new(),
216 deletes: Vec::new(),
217 };
218 encoded_bytes = base_bytes;
219 }
220 encoded_bytes += addition + usize::from(!row.puts.is_empty());
221 row.puts.push((key, value));
222 }
223 if !row.puts.is_empty() {
224 codec::encode_relay(&row)?;
225 plan.relay.push(row);
226 }
227 }
228 }
229 Ok(plan)
230}
231
232#[derive(Debug, Clone, Copy, PartialEq, Eq)]
234pub struct LocatedObject {
235 pub pack: Hash,
237 pub value: IndexValue,
239}
240
241#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
244pub enum LookupError {
245 #[error("object index row cap exceeded")]
251 TooManyRows,
252 #[error("object index page cap exceeded")]
255 TooManyPages,
256 #[error("object index membership-read cap exceeded")]
261 TooManyMembershipReads,
262}
263
264pub type ObjectLookup = Result<Option<LocatedObject>, LookupError>;
266pub type PresenceLookup = Result<bool, LookupError>;
268
269struct IdScan {
270 partition: Partition,
271 start: Key,
272 end: Key,
273 after: Option<super::Cursor>,
274 rows: Vec<(Hash, Value)>,
275 done: bool,
276 reason: Option<LookupError>,
277}
278
279#[allow(clippy::too_many_lines)] async fn scan_all<S: NamespaceStore>(
284 store: &S,
285 shards: &dyn ShardMap,
286 repo: &RepoId,
287 ids: &[Hash],
288) -> Result<BTreeMap<Hash, IdScan>, StoreError> {
289 let mut scans: BTreeMap<Hash, IdScan> = BTreeMap::new();
290 let mut order = Vec::new();
291 for id in ids {
292 if scans.contains_key(id) {
293 continue;
294 }
295 let (start, end) = keys::object_index_range(&repo.name, id);
296 scans.insert(
297 *id,
298 IdScan {
299 partition: shards.object_index(repo, id),
300 start,
301 end,
302 after: None,
303 rows: Vec::new(),
304 done: false,
305 reason: None,
306 },
307 );
308 order.push(*id);
309 }
310 let mut calls = 0;
311 let mut retained_bytes = 0;
312 let mut rotations: BTreeMap<Partition, usize> = BTreeMap::new();
313 loop {
314 let mut groups: BTreeMap<Partition, Vec<Hash>> = BTreeMap::new();
315 for id in &order {
316 let Some(scan) = scans.get_mut(id) else {
317 continue;
318 };
319 if scan.done {
320 continue;
321 }
322 if scan.rows.len() >= MAX_LOOKUP_ROWS {
323 scan.done = true;
324 scan.reason = Some(LookupError::TooManyRows);
325 continue;
326 }
327 groups.entry(scan.partition.clone()).or_default().push(*id);
328 }
329 if groups.is_empty() {
330 return Ok(scans);
331 }
332 let mut served = false;
333 for (partition, mut ids) in groups {
334 if calls == MAX_LOOKUP_PAGES {
335 for id in ids {
336 if let Some(scan) = scans.get_mut(&id) {
337 scan.done = true;
338 scan.reason = Some(LookupError::TooManyPages);
339 }
340 }
341 continue;
342 }
343 let n = ids.len();
344 ids.rotate_left(rotations.get(&partition).copied().unwrap_or(0) % n);
345 let mut remaining_rows = SCAN_CALL_ROWS;
346 let ranges: Vec<_> = ids
347 .iter()
348 .scan(&mut remaining_rows, |remaining, id| {
349 if **remaining == 0 {
350 return None;
351 }
352 let scan = &scans[id];
353 let limit = SCAN_PAGE_ROWS
354 .min(**remaining)
355 .min(u32::try_from(MAX_LOOKUP_ROWS - scan.rows.len()).unwrap_or(u32::MAX));
356 **remaining -= limit;
357 Some(RangeScan {
358 start: scan.start.clone(),
359 end: scan.end.clone(),
360 after: scan.after.clone(),
361 limit,
362 })
363 })
364 .collect();
365 let pages = store.scan_many(&partition, &ranges).await?;
366 calls += 1;
367 if pages.is_empty() || pages.len() > ranges.len() {
368 return Err(StoreError::Corrupt(
369 "invalid scan_many served prefix".into(),
370 ));
371 }
372 *rotations.entry(partition).or_default() += pages.len();
373 served = true;
374 for ((id, range), page) in ids.iter().zip(&ranges).zip(pages) {
375 if page.entries.len() > range.limit as usize {
376 return Err(StoreError::Corrupt(
377 "object index scan exceeded limit".into(),
378 ));
379 }
380 let Some(scan) = scans.get_mut(id) else {
381 return Err(StoreError::Corrupt("missing object index scan".into()));
382 };
383 for (key, value) in page.entries {
384 match keys::parse(&key) {
385 Some(ParsedKey::ObjectIndex {
386 repo: found,
387 object,
388 pack_id,
389 }) if found == repo.name && object == *id => {
390 let size = CANDIDATE_OVERHEAD.saturating_add(value.as_bytes().len());
391 if size > CANDIDATE_BYTES - retained_bytes {
392 for scan in scans.values_mut().filter(|s| !s.done) {
395 scan.done = true;
396 scan.reason = Some(LookupError::TooManyPages);
397 }
398 return Ok(scans);
399 }
400 retained_bytes += size;
401 scan.rows.push((pack_id, value));
402 }
403 _ => return Err(StoreError::Corrupt("malformed object index key".into())),
404 }
405 }
406 match page.next {
407 Some(cursor) => scan.after = Some(cursor),
408 None => scan.done = true,
409 }
410 }
411 }
412 if !served {
413 return Ok(scans);
414 }
415 }
416}
417
418pub async fn locate_many<S: NamespaceStore>(
428 store: &S,
429 shards: &dyn ShardMap,
430 repo: &RepoId,
431 ids: &[Hash],
432) -> Result<Vec<ObjectLookup>, StoreError> {
433 if ids.len() > MAX_LOOKUP_IDS {
434 return Err(StoreError::Invalid("too many object ids".into()));
435 }
436 let scans = scan_all(store, shards, repo, ids).await?;
437 let mut truncated: BTreeMap<Hash, LookupError> = scans
438 .iter()
439 .filter_map(|(id, scan)| scan.reason.map(|reason| (*id, reason)))
440 .collect();
441 let mut order: Vec<_> = scans.keys().copied().collect();
444 order.sort_by_key(|id| (scans[id].rows.len(), *id));
445 let mut packs: BTreeMap<Partition, BTreeSet<Hash>> = BTreeMap::new();
446 let mut admitted: BTreeMap<Hash, usize> = BTreeMap::new();
447 let mut membership_reads = 0;
448 for id in order {
449 let rows = &scans[&id].rows;
450 let mut count = 0;
451 for (pack, _) in rows {
452 let partition = shards.membership(repo, &BlobKey::pack(*pack));
453 let group = packs.get(&partition);
454 let new_chunk =
455 group.is_none_or(|g| !g.contains(pack) && g.len() % MEMBERSHIP_CHUNK == 0);
456 if new_chunk {
457 if membership_reads == MAX_LOOKUP_MEMBERSHIP_READS {
458 break;
459 }
460 membership_reads += 1;
461 }
462 packs.entry(partition).or_default().insert(*pack);
463 count += 1;
464 }
465 if count < rows.len() {
466 truncated
467 .entry(id)
468 .or_insert(LookupError::TooManyMembershipReads);
469 }
470 admitted.insert(id, count);
471 }
472 let mut members = BTreeSet::new();
473 for (partition, ids) in packs {
474 let ids: Vec<_> = ids.into_iter().collect();
475 for chunk in ids.chunks(MEMBERSHIP_CHUNK) {
476 let keys: Vec<_> = chunk
479 .iter()
480 .map(|pack| keys::membership(&repo.name, pack))
481 .collect();
482 let values = store.get_many(&partition, &keys).await?;
483 if values.len() != chunk.len() {
484 return Err(StoreError::Corrupt("short membership get_many".into()));
485 }
486 for (pack, value) in chunk.iter().zip(values) {
487 if value.is_some() {
488 members.insert(*pack);
489 }
490 }
491 }
492 }
493 ids.iter()
494 .map(|id| {
495 let rows = &scans[id].rows;
496 let prefix = &rows[..admitted.get(id).copied().unwrap_or(0).min(rows.len())];
497 let found = prefix
500 .iter()
501 .find(|(pack, _)| members.contains(pack))
502 .map(|(pack, value)| {
503 Ok::<LocatedObject, StoreError>(LocatedObject {
504 pack: *pack,
505 value: codec::decode_object_index(id, value)?,
506 })
507 })
508 .transpose()?;
509 if found.is_none()
510 && let Some(reason) = truncated.get(id)
511 {
512 Ok(Err(*reason))
513 } else {
514 Ok(Ok(found))
515 }
516 })
517 .collect()
518}
519
520pub async fn contains_many<S: NamespaceStore>(
522 store: &S,
523 shards: &dyn ShardMap,
524 repo: &RepoId,
525 ids: &[Hash],
526) -> Result<Vec<PresenceLookup>, StoreError> {
527 Ok(locate_many(store, shards, repo, ids)
528 .await?
529 .into_iter()
530 .map(|location| location.map(|location| location.is_some()))
531 .collect())
532}
533
534pub async fn holds_any<S: NamespaceStore>(
545 store: &S,
546 shards: &dyn ShardMap,
547 repo: &RepoId,
548 ids: &[Hash],
549) -> Result<PresenceLookup, StoreError> {
550 let answers = contains_many(store, shards, repo, ids).await?;
551 if answers.iter().any(|answer| matches!(answer, Ok(true))) {
552 return Ok(Ok(true));
553 }
554 if let Some(reason) = answers.into_iter().find_map(Result::err) {
555 return Ok(Err(reason));
556 }
557 Ok(Ok(false))
558}
559
560#[cfg(test)]
561mod tests {
562 use super::*;
563 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
564
565 use crate::pipeline::{D34Shards, SinglePartition};
566 use crate::repo::{NamespaceKey, RepoName};
567 use crate::{
568 Batch, BatchOutcome, Cursor, MemoryKv, PartitionStats, ScanPage, StoreCapabilities,
569 };
570
571 #[derive(Debug)]
572 struct EmptyPageOnce {
573 inner: MemoryKv,
574 empty_once: AtomicBool,
575 get_many_calls: AtomicUsize,
576 scan_many_calls: AtomicUsize,
577 }
578
579 impl NamespaceStore for EmptyPageOnce {
580 fn capabilities(&self) -> StoreCapabilities {
581 self.inner.capabilities()
582 }
583 async fn get(&self, p: &Partition, k: &Key) -> Result<Option<Value>, StoreError> {
584 self.inner.get(p, k).await
585 }
586 async fn get_many(
587 &self,
588 p: &Partition,
589 keys: &[Key],
590 ) -> Result<Vec<Option<Value>>, StoreError> {
591 self.get_many_calls.fetch_add(1, Ordering::SeqCst);
592 self.inner.get_many(p, keys).await
593 }
594 async fn scan(
595 &self,
596 p: &Partition,
597 start: &Key,
598 end: &Key,
599 after: Option<&Cursor>,
600 limit: u32,
601 ) -> Result<ScanPage, StoreError> {
602 if self.empty_once.swap(false, Ordering::SeqCst) {
603 let first = self.inner.scan(p, start, end, after, 1).await?;
604 return Ok(ScanPage {
605 entries: Vec::new(),
606 next: first.next,
607 });
608 }
609 self.inner.scan(p, start, end, after, limit).await
610 }
611 async fn scan_many(
612 &self,
613 p: &Partition,
614 ranges: &[RangeScan],
615 ) -> Result<Vec<ScanPage>, StoreError> {
616 self.scan_many_calls.fetch_add(1, Ordering::SeqCst);
617 let mut pages = Vec::with_capacity(ranges.len());
618 for range in ranges {
619 pages.push(
620 self.scan(
621 p,
622 &range.start,
623 &range.end,
624 range.after.as_ref(),
625 range.limit,
626 )
627 .await?,
628 );
629 }
630 Ok(pages)
631 }
632 async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
633 self.inner.apply(p, batch).await
634 }
635 async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
636 self.inner.stats(p).await
637 }
638 async fn probe(&self) -> Result<(), StoreError> {
639 self.inner.probe().await
640 }
641 }
642
643 fn repo(name: &str) -> RepoId {
644 RepoId {
645 namespace: NamespaceKey::deployment_default(),
646 name: RepoName::new(name).unwrap(),
647 }
648 }
649
650 fn raw(offset: u64) -> IndexValue {
651 IndexValue {
652 frame_offset: offset,
653 frame_length: 17,
654 wire_type: 0,
655 decoded_size: 42,
656 chain_depth: 0,
657 delta_base: None,
658 }
659 }
660
661 fn source() -> Partition {
662 Partition::Ref {
663 ns: NamespaceKey::deployment_default(),
664 repo: RepoName::new("a").unwrap(),
665 shard_ref: "refs/heads/main".into(),
666 }
667 }
668
669 #[test]
670 fn binary_value_golden_and_validation() {
671 let object = [0x12; 32];
672 let row = IndexValue {
673 frame_offset: 0x0102_0304_0506_0708,
674 frame_length: 17,
675 wire_type: 2,
676 decoded_size: 42,
677 chain_depth: 3,
678 delta_base: Some([0x33; 32]),
679 };
680 let value = codec::encode_object_index(&object, &row).unwrap();
681 let golden = [
682 &b"\x01\x01\x02\x03\x04\x05\x06\x07\x08"[..],
683 &17_u64.to_be_bytes(),
684 &[2],
685 &42_u64.to_be_bytes(),
686 &3_u32.to_be_bytes(),
687 &[1],
688 &[0x33; 32],
689 ]
690 .concat();
691 assert_eq!(value.as_bytes(), golden);
692 assert_eq!(codec::decode_object_index(&object, &value).unwrap(), row);
693 let raw_value = codec::encode_object_index(&object, &raw(0)).unwrap();
694 assert_eq!(raw_value.as_bytes().len(), 31);
695 assert_eq!(
696 codec::decode_object_index(&object, &raw_value).unwrap(),
697 raw(0)
698 );
699 for bad in [
700 Value::new(vec![]),
701 Value::new([&[2], &golden[1..]].concat()),
702 Value::new(golden[..62].to_vec()),
703 Value::new([&golden[..30], &[0], &golden[31..]].concat()),
704 ] {
705 assert!(matches!(
706 codec::decode_object_index(&object, &bad),
707 Err(StoreError::Corrupt(_))
708 ));
709 }
710 assert!(
711 codec::encode_object_index(
712 &object,
713 &IndexValue {
714 wire_type: 1,
715 ..raw(0)
716 }
717 )
718 .is_err()
719 );
720 assert!(
721 codec::encode_object_index(
722 &object,
723 &IndexValue {
724 frame_length: 0,
725 ..raw(0)
726 }
727 )
728 .is_err()
729 );
730 let deep = IndexValue {
731 chain_depth: 300,
732 ..row
733 };
734 assert_eq!(
735 codec::decode_object_index(
736 &object,
737 &codec::encode_object_index(&object, &deep).unwrap()
738 )
739 .unwrap(),
740 deep
741 );
742 for invalid in [
743 IndexValue {
744 decoded_size: MAX_RAW_OBJECT_SIZE as u64 + 1,
745 ..raw(0)
746 },
747 IndexValue {
748 frame_length: u64::from(u32::MAX) + 6,
749 ..raw(0)
750 },
751 IndexValue {
752 chain_depth: u32::from(u16::MAX) + 1,
753 ..deep
754 },
755 IndexValue {
756 delta_base: Some(object),
757 ..row
758 },
759 ] {
760 assert!(codec::encode_object_index(&object, &invalid).is_err());
761 }
762 }
763
764 #[test]
765 fn planner_chunks_deduplicates_and_is_deterministic() {
766 let r = repo("a");
767 let pack = [9; 32];
768 let entries: Vec<_> = (0..220)
769 .map(|i| {
770 let mut object = [0; 32];
771 object[30..].copy_from_slice(&u16::try_from(i).unwrap().to_be_bytes());
772 IndexEntry {
773 object,
774 value: raw(i),
775 }
776 })
777 .collect();
778 let mut duplicate = entries.clone();
779 duplicate.push(IndexEntry {
780 object: entries[0].object,
781 value: raw(999),
782 });
783 let p = Partition::Namespace(r.namespace.clone());
784 let plan = plan_index_rows(&SinglePartition, &r, &p, &pack, &duplicate, 7).unwrap();
785 assert_eq!(
786 plan,
787 plan_index_rows(&SinglePartition, &r, &p, &pack, &duplicate, 7).unwrap()
788 );
789 assert!(plan.relay.is_empty());
790 assert_eq!(
791 plan.direct
792 .iter()
793 .map(|batch| batch.puts.len())
794 .collect::<Vec<_>>(),
795 [96, 96, 28]
796 );
797 assert_eq!(
798 plan.direct[0].puts[0].1,
799 codec::encode_object_index(&entries[0].object, &raw(0)).unwrap()
800 );
801 for batch in &plan.direct {
802 let bytes: usize = batch
803 .puts
804 .iter()
805 .map(|(k, v)| k.as_bytes().len() + v.as_bytes().len())
806 .sum();
807 assert!(bytes <= MAX_BATCH_BYTES);
808 }
809 let relay = plan_index_rows(&D34Shards, &r, &source(), &pack, &duplicate, 7).unwrap();
810 assert!(relay.direct.is_empty());
811 assert_eq!(relay.relay.len(), 3);
812 assert_eq!(
813 relay
814 .relay
815 .iter()
816 .map(|row| row.puts.len())
817 .collect::<Vec<_>>(),
818 [96, 96, 28]
819 );
820 for row in relay.relay {
821 assert!(codec::encode_relay(&row).unwrap().as_bytes().len() <= MAX_VALUE_BYTES);
822 }
823 }
824
825 #[test]
826 fn planner_spans_every_prefix() {
827 let r = repo("a");
828 let entries: Vec<_> = (0..4096_u16)
829 .map(|prefix| {
830 let mut object = [0; 32];
831 object[0] = u8::try_from(prefix >> 4).unwrap();
832 object[1] = u8::try_from((prefix & 0x0f) << 4).unwrap();
833 IndexEntry {
834 object,
835 value: raw(u64::from(prefix)),
836 }
837 })
838 .collect();
839 let plan = plan_index_rows(&D34Shards, &r, &source(), &[8; 32], &entries, 7).unwrap();
840 assert_eq!(plan.relay.len(), 4096);
841 assert!(plan.relay.iter().all(|row| row.puts.len() == 1));
842 assert_eq!(
843 plan,
844 plan_index_rows(&D34Shards, &r, &source(), &[8; 32], &entries, 7).unwrap()
845 );
846 assert_eq!(
847 plan.relay
848 .iter()
849 .map(|row| &row.target)
850 .collect::<BTreeSet<_>>()
851 .len(),
852 4096
853 );
854 }
855
856 #[tokio::test]
857 async fn planner_bytes_do_not_depend_on_repository_state() {
860 let r = repo("a");
861 let object = [0x12; 32];
862 let pack = [0x34; 32];
863 let entry = IndexEntry {
864 object,
865 value: IndexValue {
866 frame_offset: 10,
867 frame_length: 20,
868 wire_type: 2,
869 decoded_size: 42,
870 chain_depth: 1,
871 delta_base: Some([0x56; 32]),
872 },
873 };
874 let source = source();
875 let empty = MemoryKv::default();
876 let member = MemoryKv::default();
877 let index_target = D34Shards.object_index(&r, &object);
878 let index_key = keys::object_index(&r.name, &object, &pack);
879 let index_value = codec::encode_object_index(&object, &entry.value).unwrap();
880 for store in [&empty, &member] {
881 store
882 .apply(
883 &index_target,
884 Batch::new().put(index_key.clone(), index_value.clone()),
885 )
886 .await
887 .unwrap();
888 }
889 member
890 .apply(
891 &D34Shards.membership(&r, &BlobKey::pack(pack)),
892 Batch::new().put(keys::membership(&r.name, &pack), Value::default()),
893 )
894 .await
895 .unwrap();
896 assert_ne!(
897 contains_many(&empty, &D34Shards, &r, &[object])
898 .await
899 .unwrap(),
900 contains_many(&member, &D34Shards, &r, &[object])
901 .await
902 .unwrap()
903 );
904 let first = plan_index_rows(&D34Shards, &r, &source, &pack, &[entry], 7).unwrap();
905 let second = plan_index_rows(&D34Shards, &r, &source, &pack, &[entry], 7).unwrap();
906 assert_eq!(first, second);
907 }
908
909 async fn membership_gate(shards: &dyn ShardMap) {
910 let store = MemoryKv::default();
911 let a = repo("a");
912 let b = repo("b");
913 let id = [0x12; 32];
914 let pack = [0x34; 32];
915 let row = codec::encode_object_index(&id, &raw(5)).unwrap();
916 let key = keys::object_index(&a.name, &id, &pack);
917 let target = shards.object_index(&a, &id);
918 assert_eq!(
919 store
920 .apply(&target, Batch::new().put(key, row))
921 .await
922 .unwrap(),
923 BatchOutcome::Committed
924 );
925 let other_target = shards.object_index(&b, &id);
926 assert_eq!(
927 store
928 .apply(
929 &other_target,
930 Batch::new().put(
931 keys::object_index(&b.name, &id, &pack),
932 codec::encode_object_index(&id, &raw(5)).unwrap(),
933 ),
934 )
935 .await
936 .unwrap(),
937 BatchOutcome::Committed
938 );
939 assert_eq!(
940 contains_many(&store, shards, &a, &[id, id]).await.unwrap(),
941 [Ok(false), Ok(false)]
942 );
943 assert_eq!(
944 contains_many(&store, shards, &b, &[id]).await.unwrap(),
945 [Ok(false)]
946 );
947 let member = keys::membership(&a.name, &pack);
948 let partition = shards.membership(&a, &BlobKey::pack(pack));
949 assert_eq!(
950 store
951 .apply(&partition, Batch::new().put(member, Value::default()))
952 .await
953 .unwrap(),
954 BatchOutcome::Committed
955 );
956 assert_eq!(
957 locate_many(&store, shards, &a, &[id]).await.unwrap(),
958 [Ok(Some(LocatedObject {
959 pack,
960 value: raw(5)
961 }))]
962 );
963 assert_eq!(
964 holds_any(&store, shards, &a, &[id]).await.unwrap(),
965 Ok(true)
966 );
967 assert_eq!(
968 holds_any(&store, shards, &b, &[id]).await.unwrap(),
969 Ok(false)
970 );
971 }
972
973 #[tokio::test]
974 async fn membership_gate_single() {
975 membership_gate(&SinglePartition).await;
976 }
977
978 #[tokio::test]
979 async fn membership_gate_d34() {
980 membership_gate(&D34Shards).await;
981 }
982
983 #[tokio::test]
984 async fn holds_any_reads_each_distinct_index_partition_once_in_first_round() {
985 let store = EmptyPageOnce {
986 inner: MemoryKv::default(),
987 empty_once: AtomicBool::new(false),
988 get_many_calls: AtomicUsize::new(0),
989 scan_many_calls: AtomicUsize::new(0),
990 };
991 let r = repo("a");
992 let mut first = [0x12; 32];
993 let mut second = first;
994 second[31] = 0x34;
995 let third = [0x34; 32];
996 first[31] = 0x56;
997 assert_eq!(
998 holds_any(&store, &D34Shards, &r, &[first, second, third])
999 .await
1000 .unwrap(),
1001 Ok(false)
1002 );
1003 assert_eq!(store.scan_many_calls.load(Ordering::SeqCst), 2);
1004 }
1005
1006 #[tokio::test]
1007 async fn lookup_pages_past_nonmember_packs() {
1008 let store = MemoryKv::default();
1009 let r = repo("a");
1010 let id = [0x12; 32];
1011 let target = D34Shards.object_index(&r, &id);
1012 for chunk in (0..130_u16).collect::<Vec<_>>().chunks(100) {
1013 let mut batch = Batch::new();
1014 for i in chunk {
1015 let mut pack = [0; 32];
1016 pack[30..].copy_from_slice(&i.to_be_bytes());
1017 batch = batch.put(
1018 keys::object_index(&r.name, &id, &pack),
1019 codec::encode_object_index(&id, &raw(u64::from(*i))).unwrap(),
1020 );
1021 }
1022 assert_eq!(
1023 store.apply(&target, batch).await.unwrap(),
1024 BatchOutcome::Committed
1025 );
1026 }
1027 let mut last = [0; 32];
1028 last[30..].copy_from_slice(&129_u16.to_be_bytes());
1029 let membership = D34Shards.membership(&r, &BlobKey::pack(last));
1030 assert_eq!(
1031 store
1032 .apply(
1033 &membership,
1034 Batch::new().put(keys::membership(&r.name, &last), Value::default())
1035 )
1036 .await
1037 .unwrap(),
1038 BatchOutcome::Committed
1039 );
1040 assert_eq!(
1041 locate_many(&store, &D34Shards, &r, &[id]).await.unwrap(),
1042 [Ok(Some(LocatedObject {
1043 pack: last,
1044 value: raw(129)
1045 }))]
1046 );
1047 }
1048
1049 #[tokio::test]
1050 async fn lookup_accepts_empty_continuation_page() {
1051 let store = EmptyPageOnce {
1052 inner: MemoryKv::default(),
1053 empty_once: AtomicBool::new(false),
1054 get_many_calls: AtomicUsize::new(0),
1055 scan_many_calls: AtomicUsize::new(0),
1056 };
1057 let r = repo("a");
1058 let id = [0x12; 32];
1059 let target = D34Shards.object_index(&r, &id);
1060 let mut member = [0; 32];
1061 member[31] = 1;
1062 for pack in [[0; 32], member] {
1063 assert_eq!(
1064 store
1065 .apply(
1066 &target,
1067 Batch::new().put(
1068 keys::object_index(&r.name, &id, &pack),
1069 codec::encode_object_index(&id, &raw(5)).unwrap()
1070 )
1071 )
1072 .await
1073 .unwrap(),
1074 BatchOutcome::Committed
1075 );
1076 }
1077 assert_eq!(
1078 store
1079 .apply(
1080 &D34Shards.membership(&r, &BlobKey::pack(member)),
1081 Batch::new().put(keys::membership(&r.name, &member), Value::default())
1082 )
1083 .await
1084 .unwrap(),
1085 BatchOutcome::Committed
1086 );
1087 store.empty_once.store(true, Ordering::SeqCst);
1088 assert_eq!(
1089 locate_many(&store, &D34Shards, &r, &[id]).await.unwrap(),
1090 [Ok(Some(LocatedObject {
1091 pack: member,
1092 value: raw(5)
1093 }))]
1094 );
1095 assert_eq!(store.get_many_calls.load(Ordering::SeqCst), 1);
1096 }
1097
1098 #[tokio::test]
1099 async fn hot_object_cap_does_not_hide_another_id() {
1100 let store = MemoryKv::default();
1101 let r = repo("a");
1102 let hot = [0x12; 32];
1103 let normal = [0x13; 32];
1104 let normal_pack = [0x55; 32];
1105 let target = D34Shards.object_index(&r, &hot);
1106 for chunk in (0..=MAX_LOOKUP_ROWS).collect::<Vec<_>>().chunks(100) {
1107 let mut batch = Batch::new();
1108 for i in chunk {
1109 let mut pack = [0; 32];
1110 pack[28..].copy_from_slice(&u32::try_from(*i).unwrap().to_be_bytes());
1111 batch = batch.put(
1112 keys::object_index(&r.name, &hot, &pack),
1113 codec::encode_object_index(&hot, &raw(5)).unwrap(),
1114 );
1115 }
1116 assert_eq!(
1117 store.apply(&target, batch).await.unwrap(),
1118 BatchOutcome::Committed
1119 );
1120 }
1121 assert_eq!(
1122 store
1123 .apply(
1124 &D34Shards.object_index(&r, &normal),
1125 Batch::new().put(
1126 keys::object_index(&r.name, &normal, &normal_pack),
1127 codec::encode_object_index(&normal, &raw(7)).unwrap(),
1128 ),
1129 )
1130 .await
1131 .unwrap(),
1132 BatchOutcome::Committed
1133 );
1134 assert_eq!(
1135 store
1136 .apply(
1137 &D34Shards.membership(&r, &BlobKey::pack(normal_pack)),
1138 Batch::new().put(keys::membership(&r.name, &normal_pack), Value::default(),),
1139 )
1140 .await
1141 .unwrap(),
1142 BatchOutcome::Committed
1143 );
1144 let normal_location = Some(LocatedObject {
1145 pack: normal_pack,
1146 value: raw(7),
1147 });
1148 assert_eq!(
1149 locate_many(&store, &D34Shards, &r, &[hot, normal])
1150 .await
1151 .unwrap(),
1152 [Err(LookupError::TooManyRows), Ok(normal_location)]
1153 );
1154 assert_eq!(
1155 contains_many(&store, &D34Shards, &r, &[hot, normal])
1156 .await
1157 .unwrap(),
1158 [Err(LookupError::TooManyRows), Ok(true)]
1159 );
1160 assert_eq!(
1161 holds_any(&store, &D34Shards, &r, &[hot, normal])
1162 .await
1163 .unwrap(),
1164 Ok(true)
1165 );
1166 assert_eq!(
1167 holds_any(&store, &D34Shards, &r, &[hot]).await.unwrap(),
1168 Err(LookupError::TooManyRows)
1169 );
1170 let first = [0; 32];
1171 assert_eq!(
1172 store
1173 .apply(
1174 &D34Shards.membership(&r, &BlobKey::pack(first)),
1175 Batch::new().put(keys::membership(&r.name, &first), Value::default()),
1176 )
1177 .await
1178 .unwrap(),
1179 BatchOutcome::Committed
1180 );
1181 assert_eq!(
1182 locate_many(&store, &D34Shards, &r, &[hot, normal])
1183 .await
1184 .unwrap(),
1185 [
1186 Ok(Some(LocatedObject {
1187 pack: first,
1188 value: raw(5)
1189 })),
1190 Ok(normal_location)
1191 ]
1192 );
1193 }
1194
1195 async fn put_rows(store: &MemoryKv, r: &RepoId, object: &Hash, packs: &[Hash]) {
1196 let target = D34Shards.object_index(r, object);
1197 for chunk in packs.chunks(100) {
1198 let mut batch = Batch::new();
1199 for pack in chunk {
1200 batch = batch.put(
1201 keys::object_index(&r.name, object, pack),
1202 codec::encode_object_index(object, &raw(5)).unwrap(),
1203 );
1204 }
1205 assert_eq!(
1206 store.apply(&target, batch).await.unwrap(),
1207 BatchOutcome::Committed
1208 );
1209 }
1210 }
1211
1212 async fn make_member(store: &MemoryKv, r: &RepoId, pack: &Hash) {
1213 assert_eq!(
1214 store
1215 .apply(
1216 &D34Shards.membership(r, &BlobKey::pack(*pack)),
1217 Batch::new().put(keys::membership(&r.name, pack), Value::default()),
1218 )
1219 .await
1220 .unwrap(),
1221 BatchOutcome::Committed
1222 );
1223 }
1224
1225 #[tokio::test]
1229 async fn spread_candidates_admit_a_pack_order_prefix() {
1230 let store = MemoryKv::default();
1231 let r = repo("a");
1232 let object = [0x21; 32];
1233 let mut packs: Vec<Hash> = (0u32..600)
1234 .map(|i| mkit_core::hash::hash(&i.to_be_bytes()))
1235 .collect();
1236 packs.sort_unstable();
1237 let partitions: BTreeSet<_> = packs
1238 .iter()
1239 .map(|pack| D34Shards.membership(&r, &BlobKey::pack(*pack)))
1240 .collect();
1241 assert!(partitions.len() > MAX_LOOKUP_MEMBERSHIP_READS);
1242 put_rows(&store, &r, &object, &packs).await;
1243 assert_eq!(
1244 locate_many(&store, &D34Shards, &r, &[object])
1245 .await
1246 .unwrap(),
1247 [Err(LookupError::TooManyMembershipReads)]
1248 );
1249 make_member(&store, &r, &packs[3]).await;
1250 assert_eq!(
1251 locate_many(&store, &D34Shards, &r, &[object])
1252 .await
1253 .unwrap(),
1254 [Ok(Some(LocatedObject {
1255 pack: packs[3],
1256 value: raw(5)
1257 }))]
1258 );
1259 }
1260
1261 #[tokio::test]
1264 async fn page_budget_round_robins_across_ids() {
1265 let store = MemoryKv::default();
1266 let r = repo("a");
1267 let mut hot = Vec::new();
1268 for h in 0u8..16 {
1269 let object = [h; 32];
1270 let packs: Vec<Hash> = (0u32..=u32::try_from(MAX_LOOKUP_ROWS).unwrap())
1271 .map(|i| {
1272 let mut pack = [0; 32];
1273 pack[28..].copy_from_slice(&i.to_be_bytes());
1274 pack
1275 })
1276 .collect();
1277 put_rows(&store, &r, &object, &packs).await;
1278 hot.push(object);
1279 }
1280 let normal = [0x77; 32];
1281 let normal_pack = [0x99; 32];
1282 put_rows(&store, &r, &normal, &[normal_pack]).await;
1283 make_member(&store, &r, &normal_pack).await;
1284 let mut ids = hot.clone();
1285 ids.push(normal);
1286 let answers = locate_many(&store, &D34Shards, &r, &ids).await.unwrap();
1287 assert_eq!(
1288 answers[16],
1289 Ok(Some(LocatedObject {
1290 pack: normal_pack,
1291 value: raw(5)
1292 }))
1293 );
1294 assert!(answers[..16].iter().all(Result::is_err));
1295 }
1296}