Skip to main content

mkit_server/store/
index.rs

1//! Repository-scoped, write-once object index rows. A row becomes visible
2//! only when its pack's repository membership row exists. Production writers
3//! are installed by WP-4.7 and WP-4.8.
4
5use 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
20/// Maximum object ids accepted by one lookup, matching the takedown named-id cap.
21pub const MAX_LOOKUP_IDS: usize = 256;
22/// Maximum index candidates read for one object id. An id that reaches this
23/// many rows with more remaining, and has no member among them, gets
24/// [`LookupError::TooManyRows`].
25pub const MAX_LOOKUP_ROWS: usize = 4096;
26/// Maximum partition-scoped scan calls in one lookup. Each call may serve a
27/// prefix of up to 256 ranges, including legal empty continuation pages.
28pub const MAX_LOOKUP_PAGES: usize = 512;
29/// At most 487 partition-scoped membership reads accompany 512 scan pages:
30/// the whole call uses at most 999 Worker subrequests.
31pub const MAX_LOOKUP_MEMBERSHIP_READS: usize = 487;
32// A 4 MiB admission budget includes payloads and a conservative 512-byte
33// charge per candidate for Vec capacity and membership BTree nodes. Together
34// with <=1,000 transient scan rows, <=256 scan states and 128-key membership
35// RPCs, a lookup retains less than 16 MiB of index state, leaving headroom in
36// the Worker's 128 MB isolate for decoding, runtime and JS transport copies.
37const 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/// The immutable location and decoded metadata of one pack entry.
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub struct IndexValue {
47    /// Offset of the complete frame in the pack.
48    pub frame_offset: u64,
49    /// Length of the complete frame on the wire.
50    pub frame_length: u64,
51    /// SPEC-PACKFILE entry type: raw, delta, zstd raw, or zstd delta.
52    pub wire_type: u8,
53    /// Reconstructed object size in bytes.
54    pub decoded_size: u64,
55    /// In-pack delta hops to the first non-in-pack-delta entry: zero for a
56    /// full object, one for an external base, otherwise one plus the base's
57    /// in-pack depth. Verification follows external bases for total depth.
58    pub chain_depth: u32,
59    /// Base id for delta entry types only.
60    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/// One verified pack entry in pack order. Repeated object ids in a pack
85/// retain the first location and value.
86#[derive(Debug, Clone, Copy, PartialEq, Eq)]
87pub struct IndexEntry {
88    /// Object id.
89    pub object: Hash,
90    /// Immutable location metadata.
91    pub value: IndexValue,
92}
93
94/// A direct batch of index upserts to the source partition.
95#[derive(Debug, Clone, PartialEq, Eq)]
96pub struct DirectIndexBatch {
97    /// The source/target partition.
98    pub target: Partition,
99    /// Key/value upserts in key order.
100    pub puts: Vec<(Key, Value)>,
101}
102
103/// Pure, deterministic plan. The caller commits each direct batch and
104/// enqueues each relay row separately from the advance batch.
105#[derive(Debug, Clone, Default, PartialEq, Eq)]
106pub struct IndexPlan {
107    /// Upserts whose target is the source partition.
108    pub direct: Vec<DirectIndexBatch>,
109    /// Upserts for other partitions, already bounded for relay delivery.
110    pub relay: Vec<RelayV1>,
111}
112
113/// Plan one pack's index rows, sorted by partition and key. The first
114/// occurrence of an object in `entries` wins. A relay row has at most
115/// [`MAX_RELAY_PUTS`] upserts and fits both value and target-batch limits.
116pub 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
127/// Plan direct idempotent upserts for every target partition. Native indexed
128/// verification uses this before committing membership; no per-object write
129/// enters the advance batch.
130pub 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                // Relay JSON renders key and value as hex strings.
196                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                    // Validation only: the enqueuer encodes the row itself.
210                    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/// A member pack and the object's immutable entry location within it.
233#[derive(Debug, Clone, Copy, PartialEq, Eq)]
234pub struct LocatedObject {
235    /// Pack whose membership makes the row visible.
236    pub pack: Hash,
237    /// Entry location and metadata.
238    pub value: IndexValue,
239}
240
241/// A bounded lookup that could not prove membership or absence for one id.
242/// Every variant fails closed for that id only; retryability differs.
243#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
244pub enum LookupError {
245    /// At least [`MAX_LOOKUP_ROWS`] rows for this object, with more
246    /// remaining, and none of those read belongs to a member pack. **Not
247    /// retryable**: it persists while the object has that many index rows in
248    /// this repository. Rows of packs that never became members are removed by
249    /// ยง13 GC (WP-5.3a); rows of member packs are not.
250    #[error("object index row cap exceeded")]
251    TooManyRows,
252    /// The call used all [`MAX_LOOKUP_PAGES`] scan pages before this id's
253    /// scan finished. **Retryable** in a call with fewer ids.
254    #[error("object index page cap exceeded")]
255    TooManyPages,
256    /// The call's membership-read budget ran out before a member was found in
257    /// this id's candidates. **Retryable** in a call with fewer ids, unless the
258    /// id's own candidates span more than [`MAX_LOOKUP_MEMBERSHIP_READS`]
259    /// membership partitions before its first member.
260    #[error("object index membership-read cap exceeded")]
261    TooManyMembershipReads,
262}
263
264/// One object's result. Backend failures still fail the whole call.
265pub type ObjectLookup = Result<Option<LocatedObject>, LookupError>;
266/// One object's membership answer.
267pub 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/// Scan candidates in rounds, with one batched call per distinct partition
280/// per round. The served-prefix cursor rotates within a partition across
281/// rounds, so a hot first id cannot indefinitely hide later ids.
282#[allow(clippy::too_many_lines)] // Keep the bounded scan state machine together.
283async 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                                // Keep only the proven pack-order prefix. The
393                                // existing page cap asks callers for smaller batches.
394                                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
418/// Locate the first member pack, in pack-id order, for each requested id.
419/// Each object has its own row cap; scan pages and membership reads have
420/// call-wide caps. A capped miss affects only its id. Each id's candidates are
421/// admitted to the membership reads as a pack-id-order prefix that fits the
422/// remaining budget, so a member early in pack order is always found.
423/// Candidate retention has a call-wide 4 MiB admission budget; reaching it
424/// uses the existing smaller-batch page-cap error. Membership keys are
425/// deduplicated and joined in 128-key chunks. Every chunk counts against the
426/// membership-read cap, including multiple chunks of one partition.
427pub 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    // Admit smaller candidate sets first, so one hot id cannot consume the
442    // membership budget needed to answer an ordinary id in the same call.
443    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            // Key and Worker base64/JSON allocations are bounded before creating
477            // any key; repository names are bounded by RepoName.
478            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            // Rows are in pack-id order and every earlier row was checked, so
498            // the first member in the admitted prefix is the first overall.
499            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
520/// Whether each object has any pack member of this repository.
521pub 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
534/// Whether this repository holds any named id, for the takedown sweep. This
535/// is a repository-index membership probe, not a `ContentIndex` hold or holder
536/// (`store/content_index.rs`), which is the global GC protection of an
537/// extracted object; the names collide. This
538/// first round performs exactly one read per distinct index partition,
539/// satisfying ยง14.3/R-133 when all requested ranges are served and each id
540/// fits one page. A served prefix or an id with more than one page needs
541/// further rounds; WP-5.6 accounts for that. A capped miss fails closed: with no hit, the first
542/// id's [`LookupError`] is returned so the caller can tell a data-dependent cap
543/// from a backend failure.
544pub 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    // The planner takes no store, so this only pins that its output ignores
858    // repository state; deriving `chain_depth` from pack bytes is WP-4.8's test.
859    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    /// Candidate packs spread over more membership partitions than the call
1226    /// budget: the first member in pack-id order is still found, and only a
1227    /// miss inside the admitted prefix reports the cap.
1228    #[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    /// Hot ids early in a request cannot spend the page budget before a later
1262    /// ordinary id gets its first page.
1263    #[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}