1use 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
51pub const INDEX_FANOUT: u16 = 4096;
54const _: () = assert!(INDEX_FANOUT == 1 << 12);
55
56pub const REF_INDEX_FANOUT: u16 = 16;
59
60pub const MAX_BLOCK_REASON_BYTES: usize = 256;
62
63pub const MAX_HOLD_TTL_MS: u64 = 24 * 60 * 60 * 1000;
67
68pub const CONTENT_APPLY_WINDOW_MS: u64 = 10_000;
72
73const MAX_ATTEMPTS: usize = 8;
76const HOLD_SCAN_PAGE: u32 = 100;
78const PRUNE_SCAN: u32 = 32;
80const GC_PRUNE_MAX: usize = 90;
83
84#[must_use]
86pub fn content_shard(object: &Hash) -> Partition {
87 Partition::ContentShard(u16::from_be_bytes([object[0], object[1]]) >> 4)
88}
89
90pub fn content_shards() -> impl Iterator<Item = Partition> {
93 (0..INDEX_FANOUT).map(Partition::ContentShard)
94}
95
96#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
98#[non_exhaustive]
99pub struct Holder {
100 pub ns: NamespaceKey,
102 pub repo: RepoName,
104}
105
106impl Holder {
107 #[must_use]
109 pub fn new(ns: NamespaceKey, repo: RepoName) -> Self {
110 Self { ns, repo }
111 }
112}
113
114#[derive(Debug, Clone, PartialEq, Eq, Default)]
116#[non_exhaustive]
117pub struct HolderPage {
118 pub holders: Vec<Holder>,
120 pub next: Option<Cursor>,
122}
123
124#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126#[non_exhaustive]
127pub struct HolderRecord {
128 pub seq: u64,
132 pub op_id: Hash,
135}
136
137impl HolderRecord {
138 #[must_use]
140 pub fn new(seq: u64, op_id: Hash) -> Self {
141 Self { seq, op_id }
142 }
143}
144
145#[derive(Debug, Clone, PartialEq, Eq)]
147#[non_exhaustive]
148pub struct BlockEntry {
149 pub reason: String,
151 pub blocked_at_ms: u64,
153}
154
155impl BlockEntry {
156 #[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
169#[non_exhaustive]
170pub struct ObjectState {
171 pub seq: u64,
173 pub changed_at_ms: u64,
175 pub holders: u64,
177 pub deleting: bool,
179}
180
181impl ObjectState {
182 #[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#[derive(Debug, Clone, PartialEq, Eq)]
196#[non_exhaustive]
197#[must_use]
198pub enum HoldOutcome {
199 Held,
201 Blocked(BlockEntry),
204}
205
206#[derive(Debug, Clone, PartialEq, Eq)]
208#[non_exhaustive]
209pub struct HolderOutcome {
210 pub newly_added: bool,
212 pub first_holder: bool,
214 pub record: HolderRecord,
216 pub blocked: Option<BlockEntry>,
220}
221
222#[derive(Debug, Clone, PartialEq, Eq)]
227#[non_exhaustive]
228pub struct GcPlan {
229 pub partition: Partition,
231 pub batch: Batch,
233}
234
235struct Seen {
237 probe: Option<Value>,
238 aux: Option<Value>,
240 blocked: Option<BlockEntry>,
241}
242
243enum 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
259fn 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
276fn 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
284fn 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
291pub(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#[derive(Debug, Clone)]
328pub struct ContentIndex<S> {
329 store: S,
330}
331
332impl<S: NamespaceStore> ContentIndex<S> {
333 pub fn new(store: S) -> Self {
335 Self { store }
336 }
337
338 pub fn store(&self) -> &S {
340 &self.store
341 }
342
343 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 const OP: Hash = [0x0f; 32];
976
977 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}