Skip to main content

mkit_server/store/
content_index.rs

1//! `ContentIndex` (PRD §5.3): the global index of what spans namespaces
2//! (holders, GC holds, the blocklist), sharded by object id over
3//! [`Partition::ContentShard`] partitions. It is a layer over any
4//! [`NamespaceStore`], not a backend trait: every per-object update is one
5//! [`Batch`] in the object's shard, which gives any backend the PRD's
6//! "atomic per-object updates".
7//!
8//! Each object has a state row (`c`, [`ObjectState`]): a change sequence,
9//! the time of the last change, the holder count and the GC `deleting`
10//! mark. Every mutation rewrites it in the same batch, guarded by `Equals`
11//! on the value it read, and prunes the expired hold rows it saw.
12//!
13//! **GC ordering (R-64).** GC deletes an object's bytes only through this
14//! sequence:
15//! 1. [`ContentIndex::collectable`] checks zero holders, zero live holds
16//!    and the grace period, and returns a [`GcPlan`]: one batch guarded on
17//!    the state row that prunes the expired holds and sets `deleting`.
18//! 2. [`ContentIndex::commit_collect`] applies it. It fails if anything
19//!    changed since step 1. Once it commits, [`ContentIndex::add_hold`] and
20//!    [`ContentIndex::add_holder`] answer a retryable
21//!    [`StoreError::Unavailable`], so no upload can dedup against bytes that
22//!    are about to go; the client retries and re-uploads.
23//! 3. GC deletes the blob (idempotent), only after step 2 committed.
24//! 4. [`ContentIndex::finish_collect`] clears `deleting`. A GC that stops
25//!    between 2 and 4 finds `deleting` set in [`ContentIndex::state`] and
26//!    resumes at step 3.
27//!
28//! **Holder rows (R-131).** A holder row is a [`HolderRecord`]: the object's
29//! post-bump change sequence and the id of the operation (the consuming
30//! ticket) that recorded it. Every holder write, including a re-record of
31//! the same holder, advances the object's sequence, and a removal is
32//! conditional on the sequence the caller read
33//! ([`ContentIndex::remove_holder`]), so a removal never undoes a holder
34//! that was re-recorded after the read (SPEC-SERVER §13.3). The count on the
35//! state row is what GC reads. The rows stay in the object's content shard.
36//!
37//! **Deadlines.** Hold and holder batches carry a `NotAfter` deadline of
38//! `now_ms + `[`CONTENT_APPLY_WINDOW_MS`]: a batch that reaches the store
39//! after its plan-time blocklist and `deleting` reads went stale writes
40//! nothing and answers a retryable [`StoreError::Unavailable`] (§14.2).
41
42use mkit_core::hash::Hash;
43
44use super::codec;
45use super::error::StoreError;
46use super::keys::{self, ParsedKey};
47use super::kv::{Batch, BatchOutcome, Cursor, Key, NamespaceStore, Precondition, Value, Write};
48use super::partition::Partition;
49use crate::repo::{NamespaceKey, RepoName};
50
51/// Object-id prefix fan-out of the content shards (and repo index shards):
52/// a fixed deployment constant, never resharded (PRD §5.3, D34).
53pub const INDEX_FANOUT: u16 = 4096;
54const _: () = assert!(INDEX_FANOUT == 1 << 12);
55
56/// Ref-name hash fan-out: a fixed deployment constant, never resharded
57/// (PRD §5.3, D34).
58pub const REF_INDEX_FANOUT: u16 = 16;
59
60/// Longest [`BlockEntry::reason`], in bytes.
61pub const MAX_BLOCK_REASON_BYTES: usize = 256;
62
63/// Longest hold, from `now_ms` to its expiry: 24 hours. A hold must
64/// outlive `MAX_APPLY_WINDOW` plus the relay-lag bound (00-plan P-21,
65/// P-23); the cap stops a caller from pinning an object indefinitely.
66pub const MAX_HOLD_TTL_MS: u64 = 24 * 60 * 60 * 1000;
67
68/// How long a hold or holder batch may take to reach the store, from the
69/// `now_ms` its plan was made at (`NotAfter`, SPEC-WRITE-GRANTS §5.5). The
70/// same bound as `MAX_APPLY_WINDOW`.
71pub const CONTENT_APPLY_WINDOW_MS: u64 = 10_000;
72
73/// Optimistic re-plans (normative rule 3) before a contended mutation
74/// fails with a retryable [`StoreError::Unavailable`].
75const MAX_ATTEMPTS: usize = 8;
76/// Page size of the hold scan in [`ContentIndex::collectable`].
77const HOLD_SCAN_PAGE: u32 = 100;
78/// Hold rows a mutation reads to prune the expired ones.
79const PRUNE_SCAN: u32 = 32;
80/// Most expired holds one GC plan deletes (the batch stays under
81/// `MAX_BATCH_OPS`); later mutations prune the rest.
82const GC_PRUNE_MAX: usize = 90;
83
84/// The content shard of `object`: its top 12 bits.
85#[must_use]
86pub fn content_shard(object: &Hash) -> Partition {
87    Partition::ContentShard(u16::from_be_bytes([object[0], object[1]]) >> 4)
88}
89
90/// Every content shard, by construction (the `Partition` enumeration
91/// rule 5).
92pub fn content_shards() -> impl Iterator<Item = Partition> {
93    (0..INDEX_FANOUT).map(Partition::ContentShard)
94}
95
96/// A repository that holds an object.
97#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
98#[non_exhaustive]
99pub struct Holder {
100    /// Namespace.
101    pub ns: NamespaceKey,
102    /// Repository.
103    pub repo: RepoName,
104}
105
106impl Holder {
107    /// A holder.
108    #[must_use]
109    pub fn new(ns: NamespaceKey, repo: RepoName) -> Self {
110        Self { ns, repo }
111    }
112}
113
114/// One page of [`ContentIndex::holders`].
115#[derive(Debug, Clone, PartialEq, Eq, Default)]
116#[non_exhaustive]
117pub struct HolderPage {
118    /// Holders in key order.
119    pub holders: Vec<Holder>,
120    /// Resume point, if more holders may follow.
121    pub next: Option<Cursor>,
122}
123
124/// A holder row's value (R-131).
125#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126#[non_exhaustive]
127pub struct HolderRecord {
128    /// The object's change sequence when the row was written (post-bump), so
129    /// the row is replaced, and a stale removal refused, by every later
130    /// write.
131    pub seq: u64,
132    /// The operation that recorded it: the consuming ticket id. Orphan
133    /// reconciliation (WP-5.3a) checks the ticket's liveness through it.
134    pub op_id: Hash,
135}
136
137impl HolderRecord {
138    /// A holder record.
139    #[must_use]
140    pub fn new(seq: u64, op_id: Hash) -> Self {
141        Self { seq, op_id }
142    }
143}
144
145/// A blocklist entry (PRD §6.7 takedown).
146#[derive(Debug, Clone, PartialEq, Eq)]
147#[non_exhaustive]
148pub struct BlockEntry {
149    /// Reason code, at most [`MAX_BLOCK_REASON_BYTES`].
150    pub reason: String,
151    /// When the object was blocked, Unix ms.
152    pub blocked_at_ms: u64,
153}
154
155impl BlockEntry {
156    /// A blocklist entry.
157    #[must_use]
158    pub fn new(reason: impl Into<String>, blocked_at_ms: u64) -> Self {
159        Self {
160            reason: reason.into(),
161            blocked_at_ms,
162        }
163    }
164}
165
166/// An object's state row. Absent means never indexed: no holders, no holds,
167/// last change at 0, not deleting.
168#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
169#[non_exhaustive]
170pub struct ObjectState {
171    /// Change sequence; every mutation increments it.
172    pub seq: u64,
173    /// Time of the last change, Unix ms (never moves backwards).
174    pub changed_at_ms: u64,
175    /// Number of holder rows.
176    pub holders: u64,
177    /// GC committed to deleting the object's bytes (see the module docs).
178    pub deleting: bool,
179}
180
181impl ObjectState {
182    /// A state row.
183    #[must_use]
184    pub fn new(seq: u64, changed_at_ms: u64, holders: u64, deleting: bool) -> Self {
185        Self {
186            seq,
187            changed_at_ms,
188            holders,
189            deleting,
190        }
191    }
192}
193
194/// The result of [`ContentIndex::add_hold`].
195#[derive(Debug, Clone, PartialEq, Eq)]
196#[non_exhaustive]
197#[must_use]
198pub enum HoldOutcome {
199    /// The hold is recorded.
200    Held,
201    /// The object is on the blocklist: nothing was written, and the upload
202    /// must be rejected.
203    Blocked(BlockEntry),
204}
205
206/// The result of [`ContentIndex::add_holder`].
207#[derive(Debug, Clone, PartialEq, Eq)]
208#[non_exhaustive]
209pub struct HolderOutcome {
210    /// The holder row was not there before.
211    pub newly_added: bool,
212    /// The object had no holder before this write.
213    pub first_holder: bool,
214    /// What the holder row now says.
215    pub record: HolderRecord,
216    /// The object's blocklist entry, if blocked: the relay then takes the
217    /// object down in the holding repo (R-75). The row is recorded either
218    /// way, so a takedown finds the holder.
219    pub blocked: Option<BlockEntry>,
220}
221
222/// A GC delete plan from [`ContentIndex::collectable`]: one batch in the
223/// object's shard, guarded on the state row it checked, that prunes the
224/// expired holds and sets `deleting`. Apply it with
225/// [`ContentIndex::commit_collect`].
226#[derive(Debug, Clone, PartialEq, Eq)]
227#[non_exhaustive]
228pub struct GcPlan {
229    /// The object's shard.
230    pub partition: Partition,
231    /// The guarded batch.
232    pub batch: Batch,
233}
234
235/// What a mutation saw, for its plan.
236struct Seen {
237    probe: Option<Value>,
238    /// The auxiliary row a mutation named (the hold a holder releases).
239    aux: Option<Value>,
240    blocked: Option<BlockEntry>,
241}
242
243/// A mutation's plan: commit writes, or stop without writing.
244enum Step<T> {
245    Commit(Vec<Write>, T),
246    Stop(T),
247    CommitBatch(Batch, T),
248}
249
250fn refuse_while_deleting(state: &ObjectState) -> Result<(), StoreError> {
251    if state.deleting {
252        return Err(StoreError::unavailable(
253            "object is being garbage-collected; retry",
254        ));
255    }
256    Ok(())
257}
258
259/// Guard the layout version `v` read: `Absent` plus a put of this binary's
260/// version, or `Equals`; refuse a newer version.
261fn guard_layout(batch: Batch, v: Option<&Value>) -> Result<Batch, StoreError> {
262    let key = keys::layout_version();
263    Ok(match v {
264        None => batch
265            .require(Precondition::Absent(key.clone()))
266            .put(key, codec::encode_u32(keys::LAYOUT_VERSION)),
267        Some(v) if codec::decode_u32(v)? > keys::LAYOUT_VERSION => {
268            return Err(StoreError::Unsupported(
269                "content shard has a newer layout version".into(),
270            ));
271        }
272        Some(v) => batch.require(Precondition::Equals(key, v.clone())),
273    })
274}
275
276/// Guard the state row read at `key`: `Equals` its value or `Absent`.
277fn guard_state(key: &Key, old: Option<&Value>) -> Precondition {
278    match old {
279        Some(v) => Precondition::Equals(key.clone(), v.clone()),
280        None => Precondition::Absent(key.clone()),
281    }
282}
283
284/// The `c` row after a change at `now_ms`.
285fn bumped(mut state: ObjectState, now_ms: u64) -> ObjectState {
286    state.seq = state.seq.wrapping_add(1);
287    state.changed_at_ms = state.changed_at_ms.max(now_ms);
288    state
289}
290
291/// A store borrowed for one call, without cloning a Worker adapter or
292/// changing its partition routing: the relay's target, and the store a
293/// [`ContentIndex`] runs over during extraction.
294pub(crate) struct BorrowedStore<'a, S>(pub(crate) &'a S);
295
296impl<S: NamespaceStore> NamespaceStore for BorrowedStore<'_, S> {
297    fn capabilities(&self) -> super::kv::StoreCapabilities {
298        self.0.capabilities()
299    }
300    async fn get(&self, p: &Partition, k: &Key) -> Result<Option<Value>, StoreError> {
301        self.0.get(p, k).await
302    }
303    async fn scan(
304        &self,
305        p: &Partition,
306        start: &Key,
307        end: &Key,
308        after: Option<&Cursor>,
309        limit: u32,
310    ) -> Result<super::kv::ScanPage, StoreError> {
311        self.0.scan(p, start, end, after, limit).await
312    }
313    async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
314        self.0.apply(p, batch).await
315    }
316    async fn stats(&self, p: &Partition) -> Result<super::kv::PartitionStats, StoreError> {
317        self.0.stats(p).await
318    }
319    async fn probe(&self) -> Result<(), StoreError> {
320        self.0.probe().await
321    }
322}
323
324/// The `ContentIndex` layer over a store's content shards. The store must
325/// accept every key class and atomic multi-key batches; otherwise every
326/// mutation fails with [`StoreError::Unsupported`].
327#[derive(Debug, Clone)]
328pub struct ContentIndex<S> {
329    store: S,
330}
331
332impl<S: NamespaceStore> ContentIndex<S> {
333    /// Wrap `store`.
334    pub fn new(store: S) -> Self {
335        Self { store }
336    }
337
338    /// The underlying store.
339    pub fn store(&self) -> &S {
340        &self.store
341    }
342
343    /// Record GC hold `hold_id` on `object` until `expires_at_ms`.
344    /// Re-adding keeps the later of the two expiries. Taken **before** the
345    /// ref-shard apply that makes the object reachable (PRD §5.3). The TTL
346    /// must exceed `MAX_APPLY_WINDOW` plus the relay-lag bound (00-plan
347    /// P-21, P-23) and is capped at [`MAX_HOLD_TTL_MS`].
348    ///
349    /// # Errors
350    /// [`StoreError::Invalid`] if the hold is already expired or longer
351    /// than [`MAX_HOLD_TTL_MS`]; a retryable [`StoreError::Unavailable`]
352    /// while GC is deleting the object.
353    pub async fn add_hold(
354        &self,
355        object: &Hash,
356        hold_id: &Hash,
357        expires_at_ms: u64,
358        now_ms: u64,
359    ) -> Result<HoldOutcome, StoreError> {
360        if expires_at_ms <= now_ms {
361            return Err(StoreError::Invalid("hold already expired".into()));
362        }
363        if expires_at_ms - now_ms > MAX_HOLD_TTL_MS {
364            return Err(StoreError::Invalid("hold exceeds MAX_HOLD_TTL_MS".into()));
365        }
366        let key = keys::hold(object, hold_id);
367        self.mutate(object, now_ms, Some(&key), None, true, |seen, state| {
368            if let Some(entry) = &seen.blocked {
369                return Ok(Step::Stop(HoldOutcome::Blocked(entry.clone())));
370            }
371            refuse_while_deleting(state)?;
372            let old = seen.probe.as_ref().map(codec::decode_hold).transpose()?;
373            let expiry = old.map_or(expires_at_ms, |old| old.max(expires_at_ms));
374            let put = Write::Put(key.clone(), codec::encode_hold(expiry));
375            Ok(Step::Commit(vec![put], HoldOutcome::Held))
376        })
377        .await
378    }
379
380    /// Extend hold `hold_id` on `object` to `expires_at_ms`, only if it is
381    /// still recorded and live at `now_ms`. A hold that is gone or expired
382    /// may already have been passed by GC, so it is not brought back: the
383    /// caller redoes from the `head` check.
384    ///
385    /// # Errors
386    /// [`StoreError::Invalid`] as for [`Self::add_hold`]; a retryable
387    /// [`StoreError::Unavailable`] while GC is deleting the object or when
388    /// the hold is no longer live.
389    pub async fn extend_hold(
390        &self,
391        object: &Hash,
392        hold_id: &Hash,
393        expires_at_ms: u64,
394        now_ms: u64,
395    ) -> Result<HoldOutcome, StoreError> {
396        if expires_at_ms <= now_ms {
397            return Err(StoreError::Invalid("hold already expired".into()));
398        }
399        if expires_at_ms - now_ms > MAX_HOLD_TTL_MS {
400            return Err(StoreError::Invalid("hold exceeds MAX_HOLD_TTL_MS".into()));
401        }
402        let key = keys::hold(object, hold_id);
403        self.mutate(object, now_ms, Some(&key), None, true, |seen, state| {
404            if let Some(entry) = &seen.blocked {
405                return Ok(Step::Stop(HoldOutcome::Blocked(entry.clone())));
406            }
407            refuse_while_deleting(state)?;
408            let old = seen.probe.as_ref().map(codec::decode_hold).transpose()?;
409            let Some(old) = old.filter(|old| *old > now_ms) else {
410                return Err(StoreError::unavailable("hold no longer live; retry"));
411            };
412            let put = Write::Put(key.clone(), codec::encode_hold(old.max(expires_at_ms)));
413            Ok(Step::Commit(vec![put], HoldOutcome::Held))
414        })
415        .await
416    }
417
418    /// Establish durable queued-holder ownership before creating a source
419    /// relay intent. Identical retries do not bump `c`; a different owner
420    /// cannot replace it. No expiration or unguarded removal is supported.
421    /// The value binds repository, source, verification job and intent.
422    pub async fn protect_pending_holder(
423        &self,
424        object: &Hash,
425        hold_id: &Hash,
426        identity: &super::PendingHolderV1,
427        now_ms: u64,
428    ) -> Result<HoldOutcome, StoreError> {
429        if identity.object != *object || identity.hold_id != *hold_id {
430            return Err(StoreError::Corrupt("pending holder key mismatch".into()));
431        }
432        let identity = identity.encode()?;
433        let key = keys::pending_holder(object, hold_id);
434        self.mutate(object, now_ms, Some(&key), None, true, |seen, state| {
435            if let Some(entry) = &seen.blocked {
436                return Ok(Step::Stop(HoldOutcome::Blocked(entry.clone())));
437            }
438            refuse_while_deleting(state)?;
439            if let Some(prior) = &seen.probe {
440                if prior != &identity {
441                    return Err(StoreError::Corrupt(
442                        "pending-holder identity changed".into(),
443                    ));
444                }
445                return Ok(Step::Stop(HoldOutcome::Held));
446            }
447            Ok(Step::Commit(
448                vec![Write::Put(key.clone(), identity.clone())],
449                HoldOutcome::Held,
450            ))
451        })
452        .await
453    }
454
455    /// Release hold `hold_id`. Normally the relay step that records the
456    /// holder row releases it in the same batch (see [`Self::add_holder`],
457    /// WP-4.10, R-75); this standalone form is for abandoned uploads.
458    pub async fn release_hold(
459        &self,
460        object: &Hash,
461        hold_id: &Hash,
462        now_ms: u64,
463    ) -> Result<(), StoreError> {
464        let key = keys::hold(object, hold_id);
465        self.mutate(object, now_ms, None, None, true, |_, _| {
466            Ok(Step::Commit(vec![Write::Delete(key.clone())], ()))
467        })
468        .await
469    }
470
471    /// Record that `holder` holds `object` (idempotent), releasing hold
472    /// `releases` in the same batch: a dedup hold is released only when its
473    /// holder row is recorded (PRD §6.7). `op_id` is the consuming ticket.
474    /// Every call advances the object's sequence and rewrites the row, also
475    /// when the holder was already recorded (SPEC-SERVER §13.3), so the
476    /// count changes only for a new holder. A blocked object is still
477    /// recorded, and reported in the outcome.
478    ///
479    /// # Errors
480    /// A retryable [`StoreError::Unavailable`] while GC is deleting the
481    /// object, or when the batch missed its `NotAfter` deadline.
482    pub async fn add_holder(
483        &self,
484        object: &Hash,
485        holder: &Holder,
486        op_id: &Hash,
487        releases: Option<&Hash>,
488        now_ms: u64,
489    ) -> Result<HolderOutcome, StoreError> {
490        self.holder_write(object, holder, op_id, releases, now_ms, false)
491            .await
492    }
493
494    /// Like [`Self::add_holder`], but a blocked object is not recorded: the
495    /// hold `releases` is deleted in the same guarded batch and the outcome
496    /// carries the block entry, so the caller fails the upload and the bytes
497    /// fall to ordinary GC (SPEC-SERVER §14.2). Used by extraction; the relay
498    /// still records blocked holders (R-75).
499    ///
500    /// # Errors
501    /// As [`Self::add_holder`].
502    pub async fn add_holder_unless_blocked(
503        &self,
504        object: &Hash,
505        holder: &Holder,
506        op_id: &Hash,
507        releases: Option<&Hash>,
508        now_ms: u64,
509    ) -> Result<HolderOutcome, StoreError> {
510        self.holder_write(object, holder, op_id, releases, now_ms, true)
511            .await
512    }
513
514    async fn holder_write(
515        &self,
516        object: &Hash,
517        holder: &Holder,
518        op_id: &Hash,
519        releases: Option<&Hash>,
520        now_ms: u64,
521        refuse_blocked: bool,
522    ) -> Result<HolderOutcome, StoreError> {
523        let key = keys::holder(object, &holder.ns, &holder.repo)?;
524        let release = releases.map(|id| keys::hold(object, id));
525        self.mutate(
526            object,
527            now_ms,
528            Some(&key),
529            release.as_ref(),
530            true,
531            |seen, state| {
532                refuse_while_deleting(state)?;
533                if refuse_blocked && let Some(entry) = &seen.blocked {
534                    let record = HolderRecord::new(bumped(*state, now_ms).seq, *op_id);
535                    let writes = release.clone().map(Write::Delete).into_iter().collect();
536                    return Ok(Step::Commit(
537                        writes,
538                        HolderOutcome {
539                            newly_added: false,
540                            first_holder: false,
541                            record,
542                            blocked: Some(entry.clone()),
543                        },
544                    ));
545                }
546                // A hold that is gone or expired no longer protects the bytes the
547                // caller relied on: GC may have passed. Retry from the `head`.
548                if release.is_some()
549                    && !seen.aux.as_ref().is_some_and(|hold| {
550                        codec::decode_hold(hold).is_ok_and(|expires| expires > now_ms)
551                    })
552                {
553                    return Err(StoreError::unavailable("hold no longer live; retry"));
554                }
555                let newly_added = seen.probe.is_none();
556                let first_holder = state.holders == 0;
557                if newly_added {
558                    state.holders += 1;
559                }
560                // The sequence this write stores: `mutate` bumps it once more.
561                let record = HolderRecord::new(bumped(*state, now_ms).seq, *op_id);
562                let mut writes = vec![Write::Put(key.clone(), codec::encode_holder(&record))];
563                writes.extend(release.clone().map(Write::Delete));
564                let outcome = HolderOutcome {
565                    newly_added,
566                    first_holder,
567                    record,
568                    blocked: seen.blocked.clone(),
569                };
570                Ok(Step::Commit(writes, outcome))
571            },
572        )
573        .await
574    }
575
576    /// The record of `holder` of `object`, if recorded.
577    pub async fn holder_record(
578        &self,
579        object: &Hash,
580        holder: &Holder,
581    ) -> Result<Option<HolderRecord>, StoreError> {
582        let key = keys::holder(object, &holder.ns, &holder.repo)?;
583        let value = self.store.get(&content_shard(object), &key).await?;
584        value.as_ref().map(codec::decode_holder).transpose()
585    }
586
587    /// Remove `holder` of `object`, only if its row still carries
588    /// `expected_seq` (the sequence of the record the caller read). `true`
589    /// if it was removed; `false` (a no-op) when the holder is absent or was
590    /// written again since (SPEC-SERVER §13.3 step 4).
591    pub async fn remove_holder(
592        &self,
593        object: &Hash,
594        holder: &Holder,
595        expected_seq: u64,
596        now_ms: u64,
597    ) -> Result<bool, StoreError> {
598        let key = keys::holder(object, &holder.ns, &holder.repo)?;
599        self.mutate(object, now_ms, Some(&key), None, true, |seen, state| {
600            let current = seen.probe.as_ref().map(codec::decode_holder).transpose()?;
601            if current.is_none_or(|record| record.seq != expected_seq) {
602                return Ok(Step::Stop(false));
603            }
604            state.holders = state.holders.saturating_sub(1);
605            Ok(Step::Commit(vec![Write::Delete(key.clone())], true))
606        })
607        .await
608    }
609
610    /// Up to `limit` holders of `object` after `after`.
611    pub async fn holders(
612        &self,
613        object: &Hash,
614        after: Option<&Cursor>,
615        limit: u32,
616    ) -> Result<HolderPage, StoreError> {
617        let (start, end) = keys::holders_of(object);
618        let page = self
619            .store
620            .scan(&content_shard(object), &start, &end, after, limit)
621            .await?;
622        let holders = page
623            .entries
624            .iter()
625            .map(|(key, _)| match keys::parse(key) {
626                Some(ParsedKey::Holder { ns, repo, .. }) => Ok(Holder { ns, repo }),
627                _ => Err(StoreError::Corrupt("malformed holder key".into())),
628            })
629            .collect::<Result<_, _>>()?;
630        Ok(HolderPage {
631            holders,
632            next: page.next,
633        })
634    }
635
636    /// Put `object` on the global blocklist (replacing any entry).
637    ///
638    /// # Errors
639    /// [`StoreError::Invalid`] if the reason exceeds
640    /// [`MAX_BLOCK_REASON_BYTES`].
641    pub async fn block(
642        &self,
643        object: &Hash,
644        entry: &BlockEntry,
645        now_ms: u64,
646    ) -> Result<(), StoreError> {
647        if entry.reason.len() > MAX_BLOCK_REASON_BYTES {
648            return Err(StoreError::Invalid("block reason too long".into()));
649        }
650        crate::takedown::directory::reserve(&self.store, object, now_ms).await?;
651        let (key, value) = (keys::block(object), codec::encode_block_entry(entry));
652        self.mutate(object, now_ms, None, None, false, |_, _| {
653            Ok(Step::Commit(
654                vec![
655                    Write::Put(key.clone(), value.clone()),
656                    Write::Put(
657                        crate::takedown::denial::legacy_descriptor_key(object),
658                        value.clone(),
659                    ),
660                ],
661                (),
662            ))
663        })
664        .await
665    }
666
667    /// Install an independent V2 action without replacing any current V1 denial.
668    pub async fn install_block_action(
669        &self,
670        object: &Hash,
671        action: &crate::takedown::denial::BlockAction,
672        now_ms: u64,
673    ) -> Result<(), StoreError> {
674        let staged =
675            crate::takedown::denial::stage_action(&self.store, object, action, now_ms).await?;
676        self.install_stored_block_action(object, &staged, now_ms)
677            .await
678    }
679
680    /// Activate already verified immutable metadata without reading source bytes.
681    pub async fn install_stored_block_action(
682        &self,
683        object: &Hash,
684        staged: &crate::takedown::denial::StoredAction,
685        now_ms: u64,
686    ) -> Result<(), StoreError> {
687        self.install_stored_block_action_with_batch(object, staged, now_ms, Batch::new())
688            .await
689    }
690
691    /// Commit source-local safety effects atomically with a new denial action.
692    pub(crate) async fn install_stored_block_action_with_batch(
693        &self,
694        object: &Hash,
695        staged: &crate::takedown::denial::StoredAction,
696        now_ms: u64,
697        effects: Batch,
698    ) -> Result<(), StoreError> {
699        use crate::takedown::denial::{action_key, decode_actions, encode_actions};
700        crate::takedown::directory::reserve(&self.store, object, now_ms).await?;
701        let key = action_key(object);
702        self.mutate(object, now_ms, Some(&key), None, false, |seen, _| {
703            let mut actions = decode_actions(seen.probe.as_ref())?;
704            match actions.binary_search_by_key(&staged.action.id, |a| a.action.id) {
705                Ok(i) if actions[i] == *staged => return Ok(Step::Stop(())),
706                Ok(_) => return Err(StoreError::Invalid("denial action identity reused".into())),
707                Err(i) => actions.insert(i, staged.clone()),
708            }
709            let value = encode_actions(actions)?;
710            let mut batch = effects.clone();
711            batch.writes.extend([
712                Write::Put(key.clone(), value.clone()),
713                Write::Put(crate::takedown::denial::descriptor_key(object), value),
714            ]);
715            Ok(Step::CommitBatch(batch, ()))
716        })
717        .await
718    }
719
720    /// Take `object` off the blocklist.
721    pub async fn unblock(&self, object: &Hash, now_ms: u64) -> Result<(), StoreError> {
722        let key = keys::block(object);
723        self.mutate(object, now_ms, None, None, false, |_, _| {
724            Ok(Step::Commit(
725                vec![
726                    Write::Delete(key.clone()),
727                    Write::Delete(crate::takedown::denial::legacy_descriptor_key(object)),
728                ],
729                (),
730            ))
731        })
732        .await
733    }
734
735    /// The blocklist entry of `object`, if blocked.
736    pub async fn blocked(&self, object: &Hash) -> Result<Option<BlockEntry>, StoreError> {
737        let rows = self
738            .store
739            .get_many(
740                &content_shard(object),
741                &[
742                    keys::block(object),
743                    crate::takedown::denial::action_key(object),
744                ],
745            )
746            .await?;
747        if rows.len() != 2 {
748            return Err(StoreError::Corrupt("short denial read".into()));
749        }
750        let legacy = rows[0]
751            .as_ref()
752            .map(codec::decode_block_entry)
753            .transpose()?;
754        let independent = crate::takedown::denial::representative(rows[1].as_ref())?;
755        Ok(legacy.or(independent))
756    }
757
758    /// The state row of `object`, if it was ever indexed.
759    pub async fn state(&self, object: &Hash) -> Result<Option<ObjectState>, StoreError> {
760        let value = self
761            .store
762            .get(&content_shard(object), &keys::object_state(object))
763            .await?;
764        value.as_ref().map(codec::decode_object_state).transpose()
765    }
766
767    /// Step 1 of the GC ordering (module docs). Whether GC may delete
768    /// `object` at `now_ms`: not already `deleting`, zero holders, zero
769    /// live holds (a hold is live while `now_ms < expires_at_ms`), and
770    /// `now_ms - last change >= grace_ms`. If so, the [`GcPlan`] that
771    /// commits the decision.
772    pub async fn collectable(
773        &self,
774        object: &Hash,
775        now_ms: u64,
776        grace_ms: u64,
777    ) -> Result<Option<GcPlan>, StoreError> {
778        let p = content_shard(object);
779        let state_key = keys::object_state(object);
780        let read = [keys::layout_version(), state_key.clone()];
781        let values = self.store.get_many(&p, &read).await?;
782        let (v, raw) = (values[0].as_ref(), values[1].as_ref());
783        let state = raw
784            .map(codec::decode_object_state)
785            .transpose()?
786            .unwrap_or_default();
787        if state.deleting
788            || state.holders > 0
789            || now_ms.saturating_sub(state.changed_at_ms) < grace_ms
790        {
791            return Ok(None);
792        }
793        // Pending holder intents remain applicable beyond ticket/hold expiry.
794        // Even unknown/corrupt rows block deletion; reconciliation must prove
795        // that no durable source intent can still deliver before removing one.
796        let (start, end) = keys::pending_holders_of(object);
797        if !self
798            .store
799            .scan(&p, &start, &end, None, 1)
800            .await?
801            .entries
802            .is_empty()
803        {
804            return Ok(None);
805        }
806        let mut expired = Vec::new();
807        let (start, end) = keys::holds_of(object);
808        let mut after = None;
809        loop {
810            let page = self
811                .store
812                .scan(&p, &start, &end, after.as_ref(), HOLD_SCAN_PAGE)
813                .await?;
814            for (key, value) in page.entries {
815                if codec::decode_hold(&value)? > now_ms {
816                    return Ok(None);
817                }
818                if expired.len() < GC_PRUNE_MAX {
819                    expired.push(Write::Delete(key));
820                }
821            }
822            match page.next {
823                Some(next) => after = Some(next),
824                None => break,
825            }
826        }
827        let mut batch = guard_layout(Batch::new(), v)?.require(guard_state(&state_key, raw));
828        batch.writes.extend(expired);
829        let mut next = bumped(state, now_ms);
830        next.deleting = true;
831        let batch = batch.put(state_key, codec::encode_object_state(&next));
832        Ok(Some(GcPlan {
833            partition: p,
834            batch,
835        }))
836    }
837
838    /// Step 2 of the GC ordering: apply `plan`. `true` if it committed and
839    /// the object is now `deleting`, so GC may delete its bytes; `false` if
840    /// anything changed since the plan was made.
841    pub async fn commit_collect(&self, plan: GcPlan) -> Result<bool, StoreError> {
842        let outcome = self.store.apply(&plan.partition, plan.batch).await?;
843        Ok(outcome == BatchOutcome::Committed)
844    }
845
846    /// Step 4 of the GC ordering: after the bytes are deleted, clear
847    /// `deleting` so the object can be uploaded again. A no-op if it is not
848    /// set.
849    pub async fn finish_collect(&self, object: &Hash, now_ms: u64) -> Result<(), StoreError> {
850        self.mutate(object, now_ms, None, None, false, |_, state| {
851            if !state.deleting {
852                return Ok(Step::Stop(()));
853            }
854            state.deleting = false;
855            Ok(Step::Commit(Vec::new(), ()))
856        })
857        .await
858    }
859
860    /// Apply one per-object mutation. Read the layout version, the state
861    /// row, the blocklist entry and `probe`, plus a page of holds; `plan`
862    /// sees them and updates the state, then either stops (nothing is
863    /// written) or returns its writes. They commit after deletes of the
864    /// expired holds read and before the new state row, guarded on the
865    /// layout version and the state row (every change to the object's rows
866    /// rewrites it). Re-plans on a lost race. With `deadline`, the batch also
867    /// requires `NotAfter(now_ms + CONTENT_APPLY_WINDOW_MS)`, and a batch
868    /// that misses it is a retryable [`StoreError::Unavailable`] (a re-plan
869    /// would reuse the same stale `now_ms`).
870    async fn mutate<T>(
871        &self,
872        object: &Hash,
873        now_ms: u64,
874        probe: Option<&Key>,
875        aux: Option<&Key>,
876        deadline: bool,
877        plan: impl Fn(&Seen, &mut ObjectState) -> Result<Step<T>, StoreError>,
878    ) -> Result<T, StoreError> {
879        let p = content_shard(object);
880        let state_key = keys::object_state(object);
881        let mut read = vec![
882            keys::layout_version(),
883            state_key.clone(),
884            keys::block(object),
885            crate::takedown::denial::action_key(object),
886        ];
887        read.extend(probe.cloned());
888        read.extend(aux.cloned());
889        let (hold_start, hold_end) = keys::holds_of(object);
890        for _ in 0..MAX_ATTEMPTS {
891            let values = self.store.get_many(&p, &read).await?;
892            if values.len() != read.len() {
893                return Err(StoreError::Corrupt("short content read".into()));
894            }
895            let mut values = values.into_iter();
896            let (v, old, blocked) = (values.next(), values.next(), values.next());
897            let (v, old, blocked) = (v.flatten(), old.flatten(), blocked.flatten());
898            let independent =
899                crate::takedown::denial::representative(values.next().flatten().as_ref())?;
900            // The optional rows follow the fixed ones, in this order.
901            let probed = probe.and_then(|_| values.next().flatten());
902            let auxiliary = aux.and_then(|_| values.next().flatten());
903            let seen = Seen {
904                probe: probed,
905                aux: auxiliary,
906                blocked: blocked
907                    .as_ref()
908                    .map(codec::decode_block_entry)
909                    .transpose()?
910                    .or(independent),
911            };
912            let holds = self
913                .store
914                .scan(&p, &hold_start, &hold_end, None, PRUNE_SCAN)
915                .await?;
916            let mut state = old
917                .as_ref()
918                .map(codec::decode_object_state)
919                .transpose()?
920                .unwrap_or_default();
921            let (mut batch, out) = match plan(&seen, &mut state)? {
922                Step::Stop(out) => return Ok(out),
923                Step::Commit(writes, out) => (
924                    Batch {
925                        preconditions: Vec::new(),
926                        writes,
927                    },
928                    out,
929                ),
930                Step::CommitBatch(batch, out) => (batch, out),
931            };
932            let writes = std::mem::take(&mut batch.writes);
933            if deadline {
934                let by = now_ms.saturating_add(CONTENT_APPLY_WINDOW_MS);
935                batch = batch.require(Precondition::NotAfter(by));
936            }
937            let mut batch =
938                guard_layout(batch, v.as_ref())?.require(guard_state(&state_key, old.as_ref()));
939            for (key, value) in holds.entries {
940                if codec::decode_hold(&value)? <= now_ms {
941                    batch = batch.delete(key);
942                }
943            }
944            batch.writes.extend(writes);
945            let batch = batch.put(
946                state_key.clone(),
947                codec::encode_object_state(&bumped(state, now_ms)),
948            );
949            match self.store.apply(&p, batch).await? {
950                BatchOutcome::Committed => return Ok(out),
951                BatchOutcome::DeadlinePassed { .. } => {
952                    return Err(StoreError::unavailable(
953                        "content index deadline passed; retry",
954                    ));
955                }
956                BatchOutcome::PreconditionFailed { .. } => {}
957            }
958        }
959        Err(StoreError::unavailable("content index update contended"))
960    }
961}
962
963#[cfg(test)]
964mod tests {
965    use std::sync::Arc;
966
967    use futures_executor::block_on;
968
969    use super::*;
970    use crate::ManualClock;
971    use crate::memory::MemoryKv;
972
973    const GRACE: u64 = 1_000;
974    /// The consuming ticket every test holder is recorded under.
975    const OP: Hash = [0x0f; 32];
976
977    /// A store whose `NotAfter` clock never reaches a test's small times.
978    fn kv() -> MemoryKv {
979        MemoryKv::with_clock(Arc::new(ManualClock::new(0)))
980    }
981
982    fn holder(repo: &str) -> Holder {
983        Holder::new(
984            NamespaceKey::deployment_default(),
985            RepoName::new(repo).unwrap(),
986        )
987    }
988
989    fn state(idx: &ContentIndex<MemoryKv>, object: &Hash) -> ObjectState {
990        block_on(idx.state(object)).unwrap().unwrap()
991    }
992
993    fn collectable(idx: &ContentIndex<MemoryKv>, object: &Hash, now: u64) -> bool {
994        block_on(idx.collectable(object, now, GRACE))
995            .unwrap()
996            .is_some()
997    }
998
999    #[test]
1000    fn expired_hold_readded_keeps_its_fresh_expiry_after_pruning() {
1001        let idx = ContentIndex::new(kv());
1002        let object = [0x53; 32];
1003        let hold = [0x35; 32];
1004        held(block_on(idx.add_hold(&object, &hold, 5, 0)));
1005        held(block_on(idx.add_hold(&object, &hold, 10_000, 6)));
1006        let key = keys::hold(&object, &hold);
1007        assert_eq!(
1008            block_on(idx.store().get(&content_shard(&object), &key)).unwrap(),
1009            Some(codec::encode_hold(10_000)),
1010            "pruning the expired snapshot must precede the fresh hold write"
1011        );
1012        assert!(!collectable(&idx, &object, GRACE + 10));
1013        held(block_on(idx.extend_hold(&object, &hold, 12_000, 7)));
1014        assert_eq!(
1015            block_on(idx.store().get(&content_shard(&object), &key)).unwrap(),
1016            Some(codec::encode_hold(12_000))
1017        );
1018        assert_eq!(state(&idx, &object).seq, 3);
1019    }
1020
1021    #[test]
1022    fn pending_holder_insert_invalidates_gc_and_retries_do_not_bump() {
1023        let idx = ContentIndex::new(kv());
1024        let object = [0x47; 32];
1025        let hold = [0x74; 32];
1026        let prior = block_on(idx.collectable(&object, GRACE, GRACE))
1027            .unwrap()
1028            .unwrap();
1029        let owner = super::super::PendingHolderV1::new(
1030            holder("repo"),
1031            Partition::Namespace(NamespaceKey::deployment_default()),
1032            OP,
1033            object,
1034            hold,
1035            [0x55; 32],
1036        )
1037        .unwrap();
1038        held(block_on(
1039            idx.protect_pending_holder(&object, &hold, &owner, GRACE),
1040        ));
1041        assert!(!block_on(idx.commit_collect(prior)).unwrap());
1042        let recorded = state(&idx, &object);
1043        held(block_on(idx.protect_pending_holder(
1044            &object,
1045            &hold,
1046            &owner,
1047            GRACE + 1,
1048        )));
1049        assert_eq!(state(&idx, &object), recorded);
1050        assert!(!collectable(&idx, &object, u64::MAX));
1051        let mut replacement = owner.clone();
1052        replacement.intent = [0x56; 32];
1053        assert!(matches!(
1054            block_on(idx.protect_pending_holder(&object, &hold, &replacement, GRACE + 1)),
1055            Err(StoreError::Corrupt(_))
1056        ));
1057        let raw = owner.encode().unwrap();
1058        assert_eq!(super::super::PendingHolderV1::decode(&raw).unwrap(), owner);
1059        let mut trailing = raw.as_bytes().to_vec();
1060        trailing.push(0);
1061        assert!(super::super::PendingHolderV1::decode(&Value::new(trailing)).is_err());
1062        assert_eq!(
1063            keys::parse(&keys::pending_holder(&object, &hold)),
1064            Some(ParsedKey::PendingHolder {
1065                object,
1066                hold_id: hold
1067            })
1068        );
1069    }
1070
1071    #[test]
1072    fn pending_holder_protection_survives_expired_hold() {
1073        let store = kv();
1074        let idx = ContentIndex::new(store);
1075        let object = [0x42; 32];
1076        let hold = [0x24; 32];
1077        held(block_on(idx.add_hold(&object, &hold, 5_000, 0)));
1078        let key = keys::pending_holder(&object, &hold);
1079        block_on(idx.store().apply(
1080            &content_shard(&object),
1081            Batch::new().put(key, Value::new(vec![1])),
1082        ))
1083        .unwrap();
1084        assert!(
1085            !collectable(&idx, &object, MAX_HOLD_TTL_MS + GRACE),
1086            "queued holder work protects bytes after the TTL expires"
1087        );
1088    }
1089
1090    /// Remove `repo`'s holder at the sequence its row carries now.
1091    fn remove(idx: &ContentIndex<MemoryKv>, object: &Hash, repo: &str, now: u64) -> bool {
1092        let Some(record) = block_on(idx.holder_record(object, &holder(repo))).unwrap() else {
1093            return block_on(idx.remove_holder(object, &holder(repo), 0, now)).unwrap();
1094        };
1095        block_on(idx.remove_holder(object, &holder(repo), record.seq, now)).unwrap()
1096    }
1097
1098    fn hold_row(idx: &ContentIndex<MemoryKv>, object: &Hash, id: &Hash) -> Option<u64> {
1099        let v = block_on(
1100            idx.store()
1101                .get(&content_shard(object), &keys::hold(object, id)),
1102        );
1103        v.unwrap().map(|v| codec::decode_hold(&v).unwrap())
1104    }
1105
1106    fn held(outcome: Result<HoldOutcome, StoreError>) {
1107        assert_eq!(outcome.unwrap(), HoldOutcome::Held);
1108    }
1109
1110    #[test]
1111    fn content_shard_is_the_top_twelve_bits() {
1112        assert_eq!(content_shard(&[0; 32]), Partition::ContentShard(0));
1113        let mut id = [0xff; 32];
1114        assert_eq!(content_shard(&id), Partition::ContentShard(4095));
1115        id[..2].copy_from_slice(&[0x12, 0x3f]);
1116        assert_eq!(content_shard(&id), Partition::ContentShard(0x123));
1117        assert_eq!(content_shards().count(), usize::from(INDEX_FANOUT));
1118    }
1119
1120    #[test]
1121    fn content_index_collectable_rules() {
1122        let idx = ContentIndex::new(kv());
1123        let obj = [7; 32];
1124        // Never indexed: nothing holds it, its last change counts as 0, and
1125        // the plan is guarded on the absent state row.
1126        assert!(!collectable(&idx, &obj, GRACE - 1));
1127        let plan = block_on(idx.collectable(&obj, GRACE, GRACE))
1128            .unwrap()
1129            .unwrap();
1130        assert!(
1131            plan.batch
1132                .preconditions
1133                .contains(&Precondition::Absent(keys::object_state(&obj)))
1134        );
1135        block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 10)).unwrap();
1136        assert!(!collectable(&idx, &obj, 10 + GRACE * 10), "held");
1137        remove(&idx, &obj, "a", 20);
1138        assert!(!collectable(&idx, &obj, 20 + GRACE - 1), "within grace");
1139        assert!(collectable(&idx, &obj, 20 + GRACE));
1140        // A live hold blocks collection; an expired one does not.
1141        held(block_on(idx.add_hold(&obj, &[1; 32], 5_000, 30)));
1142        assert!(!collectable(&idx, &obj, 30 + GRACE));
1143        assert!(!collectable(&idx, &obj, 4_999));
1144        assert!(collectable(&idx, &obj, 5_000), "a hold ends at its expiry");
1145        // A plan fails once anything changes after it was made.
1146        let stale = block_on(idx.collectable(&obj, 5_001, GRACE))
1147            .unwrap()
1148            .unwrap();
1149        let entry = BlockEntry::new("dmca", 9);
1150        block_on(idx.block(&obj, &entry, 5_002)).unwrap();
1151        assert!(!block_on(idx.commit_collect(stale)).unwrap());
1152        assert!(!state(&idx, &obj).deleting);
1153        // A blocklist entry is not a hold. The plan prunes the expired hold
1154        // and marks the object deleting.
1155        let plan = block_on(idx.collectable(&obj, 5_002 + GRACE, GRACE))
1156            .unwrap()
1157            .unwrap();
1158        assert!(block_on(idx.commit_collect(plan)).unwrap());
1159        assert!(state(&idx, &obj).deleting);
1160        assert_eq!(hold_row(&idx, &obj, &[1; 32]), None);
1161        assert!(!collectable(&idx, &obj, u64::MAX), "already deleting");
1162    }
1163
1164    #[test]
1165    fn content_index_gc_commit_beats_a_racing_upload() {
1166        let idx = ContentIndex::new(kv());
1167        let obj = [6; 32];
1168        block_on(idx.release_hold(&obj, &[0; 32], 1)).unwrap();
1169        let plan = block_on(idx.collectable(&obj, 1 + GRACE, GRACE))
1170            .unwrap()
1171            .unwrap();
1172        // GC commits first: the upload's hold must fail retryably rather
1173        // than let the upload dedup against bytes GC is about to delete.
1174        assert!(block_on(idx.commit_collect(plan)).unwrap());
1175        let now = 2 + GRACE;
1176        assert!(matches!(
1177            block_on(idx.add_hold(&obj, &[1; 32], now + 100, now)),
1178            Err(StoreError::Unavailable(_))
1179        ));
1180        assert!(matches!(
1181            block_on(idx.add_holder(&obj, &holder("a"), &OP, None, now)),
1182            Err(StoreError::Unavailable(_))
1183        ));
1184        assert_eq!(hold_row(&idx, &obj, &[1; 32]), None);
1185        // After the blob delete, GC clears the mark and uploads resume.
1186        block_on(idx.finish_collect(&obj, now)).unwrap();
1187        let s = state(&idx, &obj);
1188        block_on(idx.finish_collect(&obj, now)).unwrap();
1189        assert_eq!(state(&idx, &obj), s, "finishing twice is a no-op");
1190        held(block_on(idx.add_hold(&obj, &[1; 32], now + 100, now)));
1191        // The other order: a hold taken first makes the plan fail.
1192        let plan = block_on(idx.collectable(&obj, now + 100 + GRACE, GRACE))
1193            .unwrap()
1194            .unwrap();
1195        held(block_on(idx.add_hold(
1196            &obj,
1197            &[2; 32],
1198            now + 200 + GRACE,
1199            now + 100 + GRACE,
1200        )));
1201        assert!(!block_on(idx.commit_collect(plan)).unwrap());
1202    }
1203
1204    #[test]
1205    fn content_index_blocklist_on_add_paths() {
1206        let idx = ContentIndex::new(kv());
1207        let obj = [5; 32];
1208        let entry = BlockEntry::new("csam", 1);
1209        block_on(idx.block(&obj, &entry, 1)).unwrap();
1210        let before = state(&idx, &obj);
1211        assert_eq!(
1212            block_on(idx.add_hold(&obj, &[1; 32], 100, 2)).unwrap(),
1213            HoldOutcome::Blocked(entry.clone())
1214        );
1215        assert_eq!(state(&idx, &obj), before, "a refused hold writes nothing");
1216        assert_eq!(hold_row(&idx, &obj, &[1; 32]), None);
1217        let outcome = block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 3)).unwrap();
1218        assert_eq!((outcome.newly_added, outcome.blocked), (true, Some(entry)));
1219        block_on(idx.unblock(&obj, 4)).unwrap();
1220        held(block_on(idx.add_hold(&obj, &[1; 32], 100, 5)));
1221        let outcome = block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 6)).unwrap();
1222        assert_eq!((outcome.newly_added, outcome.blocked), (false, None));
1223    }
1224
1225    #[test]
1226    fn content_index_holds_extend_prune_and_cap() {
1227        let idx = ContentIndex::new(kv());
1228        let obj = [4; 32];
1229        held(block_on(idx.add_hold(&obj, &[1; 32], 500, 1)));
1230        held(block_on(idx.add_hold(&obj, &[1; 32], 300, 2)));
1231        assert_eq!(hold_row(&idx, &obj, &[1; 32]), Some(500), "never shortened");
1232        held(block_on(idx.add_hold(&obj, &[1; 32], 700, 3)));
1233        assert_eq!(hold_row(&idx, &obj, &[1; 32]), Some(700));
1234        held(block_on(idx.add_hold(&obj, &[2; 32], 50, 4)));
1235        // A later mutation deletes the expired hold rows it reads.
1236        block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 600)).unwrap();
1237        assert_eq!(hold_row(&idx, &obj, &[2; 32]), None);
1238        assert_eq!(hold_row(&idx, &obj, &[1; 32]), Some(700));
1239        assert!(matches!(
1240            block_on(idx.add_hold(&obj, &[3; 32], 10 + MAX_HOLD_TTL_MS + 1, 10)),
1241            Err(StoreError::Invalid(_))
1242        ));
1243        held(block_on(idx.add_hold(
1244            &obj,
1245            &[3; 32],
1246            10 + MAX_HOLD_TTL_MS,
1247            10,
1248        )));
1249    }
1250
1251    #[test]
1252    fn content_index_holder_add_idempotent() {
1253        let idx = ContentIndex::new(kv());
1254        let obj = [8; 32];
1255        for i in 0..3 {
1256            let outcome = block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 1)).unwrap();
1257            assert_eq!(outcome.newly_added, i == 0);
1258        }
1259        block_on(idx.add_holder(&obj, &holder("b"), &OP, None, 1)).unwrap();
1260        assert_eq!(state(&idx, &obj).holders, 2);
1261        let page = block_on(idx.holders(&obj, None, 1)).unwrap();
1262        assert_eq!(page.holders, vec![holder("a")]);
1263        let rest = block_on(idx.holders(&obj, page.next.as_ref(), 10)).unwrap();
1264        assert_eq!((rest.holders, rest.next), (vec![holder("b")], None));
1265        for _ in 0..2 {
1266            remove(&idx, &obj, "a", 2);
1267        }
1268        assert_eq!(state(&idx, &obj).holders, 1);
1269        // Recording a holder releases its dedup hold in the same batch.
1270        held(block_on(idx.add_hold(&obj, &[2; 32], 100, 3)));
1271        let before = state(&idx, &obj).seq;
1272        block_on(idx.add_holder(&obj, &holder("c"), &OP, Some(&[2; 32]), 4)).unwrap();
1273        assert_eq!(state(&idx, &obj).seq, before + 1);
1274        assert_eq!(state(&idx, &obj).holders, 2);
1275        assert_eq!(hold_row(&idx, &obj, &[2; 32]), None);
1276    }
1277
1278    #[test]
1279    fn content_index_every_mutation_bumps_last_change() {
1280        let idx = ContentIndex::new(kv());
1281        let obj = [9; 32];
1282        let entry = BlockEntry::new("r", 1);
1283        let mut last = ObjectState::default();
1284        for (i, now) in (1_u64..).zip([10, 20, 15, 30, 40, 50, 60, 70]) {
1285            let last_seq = last.seq;
1286            match i {
1287                1 => block_on(idx.add_hold(&obj, &[1; 32], 99, now)).map(drop),
1288                2 | 3 => block_on(idx.add_holder(&obj, &holder("a"), &OP, None, now)).map(drop),
1289                4 => block_on(idx.remove_holder(&obj, &holder("a"), last_seq, now)).map(drop),
1290                5 => block_on(idx.release_hold(&obj, &[9; 32], now)),
1291                6 => block_on(idx.block(&obj, &entry, now)),
1292                7 => block_on(idx.unblock(&obj, now)),
1293                _ => block_on(idx.release_hold(&obj, &[1; 32], now)),
1294            }
1295            .unwrap();
1296            let s = state(&idx, &obj);
1297            assert_eq!(s.seq, i, "mutation {i} bumps the sequence");
1298            // The change time never moves backwards (the 15 after 20).
1299            assert_eq!(s.changed_at_ms, now.max(last.changed_at_ms));
1300            last = s;
1301        }
1302        let v = block_on(
1303            idx.store()
1304                .get(&content_shard(&obj), &keys::layout_version()),
1305        );
1306        assert_eq!(v.unwrap(), Some(codec::encode_u32(keys::LAYOUT_VERSION)));
1307        assert!(matches!(
1308            block_on(idx.add_hold(&obj, &[1; 32], 5, 5)),
1309            Err(StoreError::Invalid(_))
1310        ));
1311        let long = BlockEntry::new("x".repeat(MAX_BLOCK_REASON_BYTES + 1), 0);
1312        assert!(matches!(
1313            block_on(idx.block(&obj, &long, 1)),
1314            Err(StoreError::Invalid(_))
1315        ));
1316        assert_eq!(state(&idx, &obj), last, "rejected calls write nothing");
1317    }
1318
1319    #[test]
1320    fn content_index_objects_in_different_shards_isolated() {
1321        let idx = ContentIndex::new(kv());
1322        let (a, mut b) = ([0x10; 32], [0x10; 32]);
1323        b[0] = 0x20;
1324        assert_ne!(content_shard(&a), content_shard(&b));
1325        block_on(idx.add_holder(&a, &holder("x"), &OP, None, 1)).unwrap();
1326        held(block_on(idx.add_hold(&a, &[1; 32], 99, 1)));
1327        let entry = BlockEntry::new("r", 1);
1328        block_on(idx.block(&a, &entry, 1)).unwrap();
1329        assert_eq!(block_on(idx.state(&b)).unwrap(), None);
1330        assert_eq!(block_on(idx.blocked(&b)).unwrap(), None);
1331        assert_eq!(block_on(idx.blocked(&a)).unwrap(), Some(entry));
1332        assert!(
1333            block_on(idx.holders(&b, None, 10))
1334                .unwrap()
1335                .holders
1336                .is_empty()
1337        );
1338        let stats = block_on(idx.store().stats(&content_shard(&b))).unwrap();
1339        assert_eq!(stats.keys, Some(0), "b's shard was never written");
1340        // Same shard, different object: rows stay apart too.
1341        let mut c = a;
1342        c[31] = 0;
1343        assert_eq!(content_shard(&a), content_shard(&c));
1344        assert!(
1345            block_on(idx.holders(&c, None, 10))
1346                .unwrap()
1347                .holders
1348                .is_empty()
1349        );
1350        assert!(collectable(&idx, &c, GRACE));
1351        assert!(!collectable(&idx, &a, GRACE * 10), "a is held");
1352    }
1353
1354    #[test]
1355    fn content_index_refuses_newer_layout_version() {
1356        let idx = ContentIndex::new(kv());
1357        let obj = [3; 32];
1358        let newer = Batch::new().put(
1359            keys::layout_version(),
1360            codec::encode_u32(keys::LAYOUT_VERSION + 1),
1361        );
1362        block_on(idx.store().apply(&content_shard(&obj), newer)).unwrap();
1363        assert!(matches!(
1364            block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 1)),
1365            Err(StoreError::Unsupported(_))
1366        ));
1367        assert!(matches!(
1368            block_on(idx.collectable(&obj, GRACE, GRACE)),
1369            Err(StoreError::Unsupported(_))
1370        ));
1371    }
1372
1373    #[test]
1374    fn content_index_holder_records_advance_seq_and_guard_removal() {
1375        let idx = ContentIndex::new(kv());
1376        let obj = [0x31; 32];
1377        let (a, b) = ([1; 32], [2; 32]);
1378        let first = block_on(idx.add_holder(&obj, &holder("a"), &a, None, 1)).unwrap();
1379        assert!(first.newly_added && first.first_holder);
1380        assert_eq!(first.record, HolderRecord::new(1, a));
1381        let second = block_on(idx.add_holder(&obj, &holder("b"), &a, None, 2)).unwrap();
1382        assert!(second.newly_added && !second.first_holder);
1383        // A re-record advances the sequence and rewrites the row, and the
1384        // count does not change.
1385        let again = block_on(idx.add_holder(&obj, &holder("a"), &b, None, 3)).unwrap();
1386        assert!(!again.newly_added && !again.first_holder);
1387        assert_eq!(again.record, HolderRecord::new(3, b));
1388        let stored = block_on(idx.holder_record(&obj, &holder("a"))).unwrap();
1389        assert_eq!(stored, Some(again.record));
1390        assert_eq!(state(&idx, &obj).holders, 2);
1391        assert_eq!(state(&idx, &obj).seq, 3);
1392        // A removal that read the first record is stale: a no-op that
1393        // writes nothing, not even a sequence bump.
1394        assert!(!block_on(idx.remove_holder(&obj, &holder("a"), 1, 4)).unwrap());
1395        assert_eq!(state(&idx, &obj).seq, 3);
1396        assert!(!block_on(idx.remove_holder(&obj, &holder("never"), 3, 4)).unwrap());
1397        assert_eq!(state(&idx, &obj).holders, 2);
1398        assert!(block_on(idx.remove_holder(&obj, &holder("a"), 3, 5)).unwrap());
1399        let last = block_on(idx.add_holder(&obj, &holder("b"), &a, None, 6)).unwrap();
1400        assert!(!last.first_holder);
1401        assert!(block_on(idx.remove_holder(&obj, &holder("b"), last.record.seq, 7)).unwrap());
1402        assert_eq!(state(&idx, &obj).holders, 0);
1403        let refill = block_on(idx.add_holder(&obj, &holder("a"), &a, None, 8)).unwrap();
1404        assert!(refill.first_holder, "no holder was left");
1405    }
1406
1407    /// A store that commits a blocklist entry for `armed` between a
1408    /// mutation's reads and its batch.
1409    struct Racing {
1410        kv: MemoryKv,
1411        armed: std::sync::Mutex<Option<(Hash, BlockEntry)>>,
1412    }
1413
1414    impl NamespaceStore for Racing {
1415        fn capabilities(&self) -> crate::store::StoreCapabilities {
1416            self.kv.capabilities()
1417        }
1418
1419        async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
1420            self.kv.get(p, key).await
1421        }
1422
1423        async fn scan(
1424            &self,
1425            p: &Partition,
1426            start: &Key,
1427            end: &Key,
1428            after: Option<&Cursor>,
1429            limit: u32,
1430        ) -> Result<crate::store::ScanPage, StoreError> {
1431            self.kv.scan(p, start, end, after, limit).await
1432        }
1433
1434        async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
1435            let race = self.armed.lock().unwrap().take();
1436            if let Some((object, entry)) = race {
1437                let raced = Batch::new()
1438                    .put(keys::block(&object), codec::encode_block_entry(&entry))
1439                    .put(
1440                        keys::object_state(&object),
1441                        codec::encode_object_state(&ObjectState::new(1, 0, 0, false)),
1442                    );
1443                self.kv.apply(&content_shard(&object), raced).await?;
1444            }
1445            self.kv.apply(p, batch).await
1446        }
1447
1448        async fn stats(&self, p: &Partition) -> Result<crate::store::PartitionStats, StoreError> {
1449            self.kv.stats(p).await
1450        }
1451
1452        async fn probe(&self) -> Result<(), StoreError> {
1453            self.kv.probe().await
1454        }
1455    }
1456
1457    #[test]
1458    fn content_index_holds_and_holders_racing_block_see_the_block() {
1459        let entry = BlockEntry::new("dmca", 1);
1460        let racing = |object: Hash| {
1461            ContentIndex::new(Racing {
1462                kv: kv(),
1463                armed: std::sync::Mutex::new(Some((object, entry.clone()))),
1464            })
1465        };
1466        // The hold's guarded batch loses to the block; its re-plan refuses.
1467        let obj = [0x41; 32];
1468        let idx = racing(obj);
1469        let hold = block_on(idx.add_hold(&obj, &[1; 32], 100, 2)).unwrap();
1470        assert_eq!(hold, HoldOutcome::Blocked(entry.clone()));
1471        assert_eq!(hold_row_in(&idx, &obj, &[1; 32]), None);
1472        // The holder is recorded anyway (a takedown must find it) and
1473        // reports the block.
1474        let obj = [0x42; 32];
1475        let idx = racing(obj);
1476        let outcome = block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 2)).unwrap();
1477        assert_eq!(outcome.blocked, Some(entry.clone()));
1478        assert!(outcome.newly_added);
1479    }
1480
1481    fn hold_row_in(idx: &ContentIndex<Racing>, object: &Hash, id: &Hash) -> Option<u64> {
1482        let key = keys::hold(object, id);
1483        let v = block_on(idx.store().get(&content_shard(object), &key)).unwrap();
1484        v.map(|v| codec::decode_hold(&v).unwrap())
1485    }
1486
1487    #[test]
1488    fn content_index_hold_and_holder_batches_carry_a_deadline() {
1489        let clock = Arc::new(ManualClock::new(0));
1490        let idx = ContentIndex::new(MemoryKv::with_clock(clock.clone()));
1491        let obj = [0x51; 32];
1492        held(block_on(idx.add_hold(&obj, &[1; 32], 100, 10)));
1493        // The store's clock is past the plan's `now + window`: nothing is
1494        // written, and the failure is the retryable kind.
1495        clock.set(i64::try_from(10 + CONTENT_APPLY_WINDOW_MS + 1).unwrap());
1496        let before = state(&idx, &obj);
1497        let hold = block_on(idx.add_hold(&obj, &[2; 32], 200, 10));
1498        assert!(matches!(hold, Err(StoreError::Unavailable(_))));
1499        let holder_added = block_on(idx.add_holder(&obj, &holder("a"), &OP, None, 10));
1500        assert!(matches!(holder_added, Err(StoreError::Unavailable(_))));
1501        let removed = block_on(idx.remove_holder(&obj, &holder("a"), 1, 10));
1502        assert!(
1503            !removed.unwrap(),
1504            "an absent holder is a no-op, before any write"
1505        );
1506        assert_eq!(state(&idx, &obj), before);
1507        assert_eq!(hold_row(&idx, &obj, &[2; 32]), None);
1508        assert_eq!(
1509            block_on(idx.holder_record(&obj, &holder("a"))).unwrap(),
1510            None
1511        );
1512        // A fresh plan time succeeds.
1513        let now = 10 + CONTENT_APPLY_WINDOW_MS + 1;
1514        assert!(block_on(idx.add_holder(&obj, &holder("a"), &OP, None, now)).is_ok());
1515    }
1516
1517    #[test]
1518    fn content_index_holder_releases_only_a_live_hold() {
1519        let idx = ContentIndex::new(kv());
1520        let obj = [0x61; 32];
1521        // No hold at all, then an expired one: the holder batch refuses and
1522        // writes nothing.
1523        let none = block_on(idx.add_holder(&obj, &holder("a"), &OP, Some(&[1; 32]), 1));
1524        assert!(matches!(none, Err(StoreError::Unavailable(_))));
1525        held(block_on(idx.add_hold(&obj, &[1; 32], 50, 2)));
1526        let before = state(&idx, &obj);
1527        let late = block_on(idx.add_holder(&obj, &holder("a"), &OP, Some(&[1; 32]), 60));
1528        assert!(matches!(late, Err(StoreError::Unavailable(_))));
1529        assert_eq!(state(&idx, &obj).holders, before.holders);
1530        assert_eq!(
1531            block_on(idx.holder_record(&obj, &holder("a"))).unwrap(),
1532            None
1533        );
1534        // A live hold is released with the holder.
1535        held(block_on(idx.add_hold(&obj, &[2; 32], 500, 70)));
1536        block_on(idx.add_holder(&obj, &holder("a"), &OP, Some(&[2; 32]), 80)).unwrap();
1537        assert_eq!(hold_row(&idx, &obj, &[2; 32]), None);
1538        assert!(
1539            block_on(idx.holder_record(&obj, &holder("a")))
1540                .unwrap()
1541                .is_some()
1542        );
1543    }
1544}