1use crate::indexed::{
3 IndexedConfig,
4 budget::{Budgeted, SliceBudget},
5};
6use crate::pipeline::ShardMap;
7use crate::store::{
8 BlockEntry, BorrowedStore, ContentIndex, Key, NamespaceStore, StoreError, Value, keys,
9};
10use crate::{RepoId, ServerError};
11use mkit_core::hash::Hash;
12use serde::{Deserialize, Serialize};
13use std::collections::BTreeSet;
14
15#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
17#[serde(deny_unknown_fields)]
18pub struct BlockAction {
19 pub id: Hash,
20 pub takedown_id: Hash,
21 pub reason: String,
22 pub blocked_at_ms: u64,
23 pub chunk_ids: Vec<Hash>,
24}
25#[derive(Serialize, Deserialize)]
26#[serde(deny_unknown_fields)]
27struct ActionsV2 {
28 version: u8,
29 actions: Vec<StoredAction>,
30}
31#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
32#[serde(deny_unknown_fields)]
33pub struct StoredAction {
34 pub action: BlockAction,
35 pub sorted_pages: bool,
36 pub chunk_count: u64,
37 pub chunk_digest: Hash,
38 pub pages: Vec<ChunkPage>,
39 pub page_owner: Hash,
40 pub page_action: Hash,
41 pub pack_scope: Option<Hash>,
42 pub pack_digest: Option<Hash>,
43}
44#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
45#[serde(deny_unknown_fields)]
46pub struct ChunkPage {
47 pub first: Hash,
48 pub last: Hash,
49 pub count: u32,
50 pub digest: Hash,
51}
52const INDEX_PREFIX: &[u8] = b"b\0\xffdenial-action-descriptors-v2-index\0";
53#[must_use]
54pub fn descriptor_key(object: &Hash) -> Key {
55 let mut bytes = INDEX_PREFIX.to_vec();
56 bytes.extend_from_slice(object);
57 Key::new(bytes)
58}
59const PAGE_HASHES: usize = crate::store::MAX_VALUE_BYTES / 32;
60#[cfg(test)]
61static STAGED_PAGE_PEAK: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
62pub(super) fn chunk_page_key(object: &Hash, action: &Hash, page: u32) -> Key {
63 let mut bytes = keys::block(object).as_bytes().to_vec();
64 bytes.extend_from_slice(b"\0chunks\0");
65 bytes.extend_from_slice(action);
66 bytes.extend_from_slice(&page.to_be_bytes());
67 Key::new(bytes)
68}
69pub async fn stage_action<S: NamespaceStore>(
72 store: &S,
73 object: &Hash,
74 action: &BlockAction,
75 now: u64,
76) -> Result<StoredAction, StoreError> {
77 if action.chunk_ids.windows(2).any(|w| w[0] >= w[1]) {
78 return Err(corrupt());
79 }
80 stage_references(store, object, action, &action.chunk_ids, true, now).await
81}
82pub(super) async fn stage_references<S: NamespaceStore>(
84 store: &S,
85 object: &Hash,
86 action: &BlockAction,
87 ids: &[Hash],
88 sorted_pages: bool,
89 now: u64,
90) -> Result<StoredAction, StoreError> {
91 let mut digest = blake3::Hasher::new();
92 let mut pages = Vec::new();
93 for (page, chunk_ids) in ids.chunks(PAGE_HASHES).enumerate() {
94 let mut chunks = chunk_ids.to_vec();
95 chunks.sort_unstable();
96 chunks.dedup();
97 let bytes = chunks.concat();
98 #[cfg(test)]
99 STAGED_PAGE_PEAK.fetch_max(
100 chunks.capacity() * 32 + bytes.capacity(),
101 std::sync::atomic::Ordering::Relaxed,
102 );
103 digest.update(&bytes);
104 let key = chunk_page_key(
105 object,
106 &action.id,
107 u32::try_from(page).map_err(|_| corrupt())?,
108 );
109 pages.push(ChunkPage {
110 first: chunks[0],
111 last: chunks[chunks.len() - 1],
112 count: u32::try_from(chunks.len()).map_err(|_| corrupt())?,
113 digest: mkit_core::hash::hash(&bytes),
114 });
115 let value = Value::new(bytes);
116 let partition = crate::store::content_shard(object);
117 let old = store.get(&partition, &key).await?;
118 if let Some(old) = old {
119 if old != value {
120 return Err(StoreError::Invalid("denial action identity reused".into()));
121 }
122 } else {
123 let batch = crate::Batch::new()
124 .require(crate::Precondition::Absent(key.clone()))
125 .require(crate::Precondition::NotAfter(
126 now.saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS),
127 ))
128 .put(key.clone(), value.clone());
129 match store.apply(&partition, batch).await? {
130 crate::BatchOutcome::Committed => {}
131 _ if store.get(&partition, &key).await?.as_ref() == Some(&value) => {}
132 _ => return Err(StoreError::unavailable("denial chunk page race")),
133 }
134 }
135 }
136 let header = BlockAction {
137 id: action.id,
138 takedown_id: action.takedown_id,
139 reason: action.reason.clone(),
140 blocked_at_ms: action.blocked_at_ms,
141 chunk_ids: Vec::new(),
142 };
143 Ok(StoredAction {
144 action: header,
145 sorted_pages,
146 chunk_count: pages.iter().map(|p| u64::from(p.count)).sum(),
147 chunk_digest: *digest.finalize().as_bytes(),
148 pages,
149 page_owner: *object,
150 page_action: action.id,
151 pack_scope: None,
152 pack_digest: None,
153 })
154}
155pub(crate) async fn page<S: NamespaceStore>(
156 store: &S,
157 row: &StoredAction,
158 index: usize,
159) -> Result<Vec<Hash>, ServerError> {
160 let desc = row.pages.get(index).ok_or_else(unavailable)?;
161 let raw = store
162 .get(
163 &crate::store::content_shard(&row.page_owner),
164 &chunk_page_key(
165 &row.page_owner,
166 &row.page_action,
167 u32::try_from(index).map_err(|_| unavailable())?,
168 ),
169 )
170 .await
171 .map_err(|_| unavailable())?
172 .ok_or_else(unavailable)?;
173 if raw.as_bytes().len() != desc.count as usize * 32
174 || mkit_core::hash::hash(raw.as_bytes()) != desc.digest
175 {
176 return Err(unavailable());
177 }
178 let ids: Vec<Hash> = raw
179 .as_bytes()
180 .chunks_exact(32)
181 .map(|bytes| bytes.try_into().map_err(|_| unavailable()))
182 .collect::<Result<_, _>>()?;
183 if ids.first() != Some(&desc.first)
184 || ids.last() != Some(&desc.last)
185 || ids.windows(2).any(|w| w[0] >= w[1])
186 {
187 return Err(unavailable());
188 }
189 Ok(ids)
190}
191pub(super) async fn chunks_intersect<S: NamespaceStore>(
192 store: &S,
193 _object: &Hash,
194 row: &StoredAction,
195 ids: &BTreeSet<Hash>,
196) -> Result<bool, ServerError> {
197 for (index, desc) in row.pages.iter().enumerate() {
198 if ids.range(desc.first..=desc.last).next().is_none() {
199 continue;
200 }
201 if page(store, row, index)
202 .await?
203 .iter()
204 .any(|id| ids.contains(id))
205 {
206 return Ok(true);
207 }
208 }
209 Ok(false)
210}
211
212#[must_use]
214pub fn action_key(id: &Hash) -> Key {
215 let mut bytes = keys::block(id).as_bytes().to_vec();
216 bytes.extend_from_slice(b"\0actions");
217 Key::new(bytes)
218}
219pub fn decode_actions(raw: Option<&Value>) -> Result<Vec<StoredAction>, StoreError> {
220 let Some(raw) = raw else {
221 return Ok(Vec::new());
222 };
223 let dto: ActionsV2 = serde_json::from_slice(raw.as_bytes()).map_err(|_| corrupt())?;
224 if dto.version != 2
225 || dto.actions.is_empty()
226 || dto
227 .actions
228 .windows(2)
229 .any(|w| w[0].action.id >= w[1].action.id)
230 || dto.actions.iter().any(|a| {
231 a.action.reason.is_empty()
232 || a.action.reason.len() > 256
233 || !a.action.chunk_ids.is_empty()
234 || a.pack_scope.is_some() != a.pack_digest.is_some()
235 || a.pages
236 .iter()
237 .any(|p| p.count == 0 || p.count as usize > PAGE_HASHES || p.first > p.last)
238 || (a.sorted_pages && a.pages.windows(2).any(|w| w[0].last >= w[1].first))
239 || a.pages.iter().map(|p| u64::from(p.count)).sum::<u64>() != a.chunk_count
240 })
241 {
242 return Err(corrupt());
243 }
244 Ok(dto.actions)
245}
246pub fn encode_actions(actions: Vec<StoredAction>) -> Result<Value, StoreError> {
247 let bytes = serde_json::to_vec(&ActionsV2 {
248 version: 2,
249 actions,
250 })
251 .map_err(|_| corrupt())?;
252 if bytes.len() > crate::store::MAX_VALUE_BYTES {
253 return Err(StoreError::Full);
254 }
255 let value = Value::new(bytes);
256 decode_actions(Some(&value))?;
257 Ok(value)
258}
259fn corrupt() -> StoreError {
260 StoreError::Corrupt("bad denial actions".into())
261}
262pub fn representative(raw: Option<&Value>) -> Result<Option<BlockEntry>, StoreError> {
263 Ok(decode_actions(raw)?
264 .first()
265 .map(|a| BlockEntry::new(&a.action.reason, a.action.blocked_at_ms)))
266}
267pub async fn denied<S: NamespaceStore>(store: &S, id: &Hash) -> Result<bool, ServerError> {
268 ContentIndex::new(BorrowedStore(store))
269 .blocked(id)
270 .await
271 .map(|entry| entry.is_some())
272 .map_err(|_| unavailable())
273}
274pub async fn require_clear<S: NamespaceStore>(store: &S, id: &Hash) -> Result<(), ServerError> {
275 if denied(store, id).await? {
276 return Err(blocked());
277 }
278 Ok(())
279}
280fn blocked() -> ServerError {
281 ServerError::permission_denied("object blocked")
282}
283fn unavailable() -> ServerError {
284 ServerError::unavailable("object storage request failed")
285}
286
287const MAX_PROOF_CONTEXT_BYTES: usize = 4 * 1024 * 1024;
289struct Target<'a> {
290 ids: &'a BTreeSet<Hash>,
291 manifests: &'a [StoredAction],
292 packs: Vec<Hash>,
293}
294impl Target<'_> {
295 async fn intersects<S: NamespaceStore>(
296 &self,
297 store: &S,
298 ids: &[Hash],
299 ) -> Result<bool, ServerError> {
300 if ids.iter().any(|id| self.ids.contains(id)) {
301 return Ok(true);
302 }
303 let requested: BTreeSet<Hash> = ids.iter().copied().collect();
304 for manifest in self.manifests {
305 if chunks_intersect(store, &manifest.page_owner, manifest, &requested).await? {
306 return Ok(true);
307 }
308 }
309 for pack in &self.packs {
310 let keys: Vec<Key> = ids
311 .iter()
312 .map(|id| super::inventory::marker_key(pack, id))
313 .collect();
314 let rows = store
315 .get_many(&crate::store::content_shard(pack), &keys)
316 .await
317 .map_err(|_| unavailable())?;
318 if rows.len() != keys.len() {
319 return Err(unavailable());
320 }
321 if rows.iter().flatten().any(|row| row.as_bytes().len() != 32) {
322 return Err(unavailable());
323 }
324 if rows.iter().any(Option::is_some) {
325 return Ok(true);
326 }
327 if super::inventory::visit(store, pack, true, |_, row| {
328 let requested = &requested;
329 async move {
330 if row.kind != 5 {
331 return Ok(false);
332 }
333 chunks_intersect(store, pack, &row.references, requested)
334 .await
335 .map_err(|_| corrupt())
336 }
337 })
338 .await
339 .map_err(|_| unavailable())?
340 {
341 return Ok(true);
342 }
343 }
344 Ok(false)
345 }
346 async fn chunks<S: NamespaceStore>(
347 &self,
348 store: &S,
349 row: &StoredAction,
350 ) -> Result<bool, ServerError> {
351 for (n, desc) in row.pages.iter().enumerate() {
352 if self.packs.is_empty()
353 && self.ids.range(desc.first..=desc.last).next().is_none()
354 && !self
355 .manifests
356 .iter()
357 .flat_map(|m| &m.pages)
358 .any(|p| p.first <= desc.last && desc.first <= p.last)
359 {
360 continue;
361 }
362 for ids in page(store, row, n).await?.chunks(256) {
363 if self.intersects(store, ids).await? {
364 return Ok(true);
365 }
366 }
367 }
368 Ok(false)
369 }
370}
371#[must_use]
372pub fn legacy_descriptor_key(object: &Hash) -> Key {
373 let mut bytes = descriptor_key(object).as_bytes().to_vec();
374 bytes.push(0);
375 Key::new(bytes)
376}
377async fn held<S: NamespaceStore>(
378 store: &S,
379 shards: &dyn ShardMap,
380 repo: &RepoId,
381 id: &Hash,
382) -> Result<bool, ServerError> {
383 crate::store::index::holds_any(store, shards, repo, &[*id])
384 .await
385 .map_err(|_| unavailable())?
386 .map_err(|_| unavailable())
387}
388async fn prove<S: NamespaceStore>(
391 store: &S,
392 shards: &dyn ShardMap,
393 repo: &RepoId,
394 target: &Target<'_>,
395) -> Result<(), ServerError> {
396 let start = Key::new(INDEX_PREFIX.to_vec());
397 let mut end = INDEX_PREFIX.to_vec();
398 *end.last_mut().ok_or_else(unavailable)? = 1;
399 let end = Key::new(end);
400 prove_shards(
401 store,
402 shards,
403 repo,
404 target,
405 &start,
406 &end,
407 WRITE_PROOF_CONCURRENCY,
408 )
409 .await?;
410 for id in target.ids {
411 require_clear(store, id).await?;
412 }
413 Ok(())
414}
415const SCANNER_PROOF_CONCURRENCY: usize = 6;
418const WRITE_PROOF_CONCURRENCY: usize = 4;
421const _: () = assert!(SCANNER_PROOF_CONCURRENCY * crate::store::MAX_VALUE_BYTES <= 4 * 1024 * 1024);
424
425async fn prove_shards<S: NamespaceStore>(
426 store: &S,
427 shards: &dyn ShardMap,
428 repo: &RepoId,
429 target: &Target<'_>,
430 start: &Key,
431 end: &Key,
432 concurrency: usize,
433) -> Result<(), ServerError> {
434 let _ = (start, end);
435 let mut directory = super::directory::Walk::start(store, concurrency)
436 .await
437 .map_err(|_| unavailable())?;
438 while let Some(id) = directory.next(store).await.map_err(|_| unavailable())? {
439 for (key, raw) in super::directory::descriptors(store, &id)
440 .await
441 .map_err(|_| unavailable())?
442 {
443 check_descriptor(store, shards, repo, target, &key, &raw).await?;
444 }
445 }
446 Ok(())
447}
448async fn check_descriptor<S: NamespaceStore>(
449 store: &S,
450 shards: &dyn ShardMap,
451 repo: &RepoId,
452 target: &Target<'_>,
453 key: &Key,
454 raw: &Value,
455) -> Result<(), ServerError> {
456 let body = key
457 .as_bytes()
458 .strip_prefix(INDEX_PREFIX)
459 .ok_or_else(unavailable)?;
460 if body.len() != 32 && !(body.len() == 33 && body[32] == 0) {
461 return Err(unavailable());
462 }
463 let object: Hash = body[..32].try_into().map_err(|_| unavailable())?;
464 if target.intersects(store, &[object]).await? {
465 return Err(blocked());
466 }
467 if body.len() == 33 {
468 crate::store::codec::decode_block_entry(raw).map_err(|_| unavailable())?;
469 if held(store, shards, repo, &object).await? {
470 let (_, row) = super::inventory::member(store, shards, repo, &object)
471 .await
472 .map_err(|_| unavailable())?;
473 if row.kind == 5 && target.chunks(store, &row.references).await? {
474 return Err(blocked());
475 }
476 }
477 return Ok(());
478 }
479 for action in decode_actions(Some(raw)).map_err(|_| unavailable())? {
480 if let Some(pack) = action.pack_scope {
481 if super::inventory::seal(store, &pack)
482 .await
483 .map_err(|_| unavailable())?
484 != action.pack_digest.ok_or_else(unavailable)?
485 {
486 return Err(unavailable());
487 }
488 let hit = super::inventory::visit(store, &pack, false, |id, row| async move {
489 if target
490 .intersects(store, &[id])
491 .await
492 .map_err(|_| corrupt())?
493 && super::inventory::is_file(store, &pack, &id).await?
494 {
495 return Ok(true);
496 }
497 if row.kind == 5
498 && held(store, shards, repo, &id)
499 .await
500 .map_err(|_| corrupt())?
501 && target
502 .chunks(store, &row.references)
503 .await
504 .map_err(|_| corrupt())?
505 {
506 return Ok(true);
507 }
508 Ok(false)
509 })
510 .await
511 .map_err(|_| unavailable())?;
512 if hit {
513 return Err(blocked());
514 }
515 } else if action.chunk_count > 0
516 && held(store, shards, repo, &object).await?
517 && target.chunks(store, &action).await?
518 {
519 return Err(blocked());
520 }
521 }
522 Ok(())
523}
524pub async fn require_repo_clear<S: NamespaceStore>(
525 store: &S,
526 shards: &dyn ShardMap,
527 repo: &RepoId,
528 ids: &BTreeSet<Hash>,
529 budget: &SliceBudget,
530) -> Result<(), ServerError> {
531 prove(
532 &Budgeted::new(store, budget),
533 shards,
534 repo,
535 &Target {
536 ids,
537 manifests: &[],
538 packs: Vec::new(),
539 },
540 )
541 .await
542}
543pub async fn require_repo_clear_budgeted<S: NamespaceStore>(
544 store: &S,
545 shards: &dyn ShardMap,
546 repo: &RepoId,
547 ids: &BTreeSet<Hash>,
548 packs: &[Hash],
549 budget: &SliceBudget,
550) -> Result<(), ServerError> {
551 let remote = Budgeted::new(store, budget);
552 for pack in packs {
553 super::inventory::visit(&remote, pack, false, |_, _| async { Ok(false) })
554 .await
555 .map_err(|_| unavailable())?;
556 }
557 prove(
558 &remote,
559 shards,
560 repo,
561 &Target {
562 ids,
563 manifests: &[],
564 packs: packs.to_vec(),
565 },
566 )
567 .await
568}
569pub async fn require_object_clear<S: NamespaceStore>(
570 store: &S,
571 shards: &dyn ShardMap,
572 repo: &RepoId,
573 id: &Hash,
574 cfg: &IndexedConfig,
575 budget: &SliceBudget,
576) -> Result<(), ServerError> {
577 let remote = Budgeted::new(store, budget);
578 let context = object_context(&remote, shards, repo, id, cfg).await?;
579 prove(&remote, shards, repo, &context.target()).await
580}
581struct ObjectContext {
582 ids: BTreeSet<Hash>,
583 manifests: Vec<StoredAction>,
584}
585impl ObjectContext {
586 fn bytes(&self) -> usize {
587 let mut bytes = self.ids.len() * 128;
588 for m in &self.manifests {
589 bytes += std::mem::size_of::<StoredAction>()
590 + m.action.reason.capacity()
591 + m.pages.capacity() * std::mem::size_of::<ChunkPage>();
592 }
593 bytes
594 }
595 fn target(&self) -> Target<'_> {
596 Target {
597 ids: &self.ids,
598 manifests: &self.manifests,
599 packs: Vec::new(),
600 }
601 }
602}
603async fn object_context<S: NamespaceStore>(
604 store: &S,
605 shards: &dyn ShardMap,
606 repo: &RepoId,
607 id: &Hash,
608 cfg: &IndexedConfig,
609) -> Result<ObjectContext, ServerError> {
610 let mut context = ObjectContext {
611 ids: BTreeSet::from([*id]),
612 manifests: Vec::new(),
613 };
614 let mut cursor = Some(*id);
615 for _ in 0..=cfg.max_delta_chain_depth {
616 let Some(id) = cursor else { break };
617 let (pack, row) = super::inventory::member(store, shards, repo, &id)
618 .await
619 .map_err(|_| unavailable())?;
620 context.ids.insert(pack);
621 if row.kind == 5 {
622 context.manifests.push(row.references);
623 if context.bytes() > MAX_PROOF_CONTEXT_BYTES {
624 return Err(unavailable());
625 }
626 }
627 cursor = row.base;
628 if let Some(base) = cursor
629 && !context.ids.insert(base)
630 {
631 return Err(unavailable());
632 }
633 }
634 if cursor.is_some() {
635 return Err(unavailable());
636 }
637 Ok(context)
638}
639#[cfg(feature = "http-objects")]
640pub(crate) async fn object_denials<S: NamespaceStore>(
641 store: &S,
642 shards: &dyn ShardMap,
643 repo: &RepoId,
644 requested: &BTreeSet<Hash>,
645 cfg: &IndexedConfig,
646) -> Result<BTreeSet<Hash>, ServerError> {
647 let mut contexts = Vec::new();
648 let mut bytes = 0;
649 for id in requested {
650 let context = object_context(store, shards, repo, id, cfg).await?;
651 bytes += context.bytes();
652 if bytes > MAX_PROOF_CONTEXT_BYTES {
653 return Err(unavailable());
654 }
655 contexts.push((*id, context));
656 }
657 let start = Key::new(INDEX_PREFIX.to_vec());
658 let mut end = INDEX_PREFIX.to_vec();
659 *end.last_mut().ok_or_else(unavailable)? = 1;
660 let end = Key::new(end);
661 let mut denied = BTreeSet::new();
662 let _ = (start, end);
663 let mut directory = super::directory::Walk::start(store, WRITE_PROOF_CONCURRENCY)
664 .await
665 .map_err(|_| unavailable())?;
666 while let Some(object) = directory.next(store).await.map_err(|_| unavailable())? {
667 for (key, raw) in super::directory::descriptors(store, &object)
668 .await
669 .map_err(|_| unavailable())?
670 {
671 for (id, context) in &contexts {
672 if denied.contains(id) {
673 continue;
674 }
675 if let Err(error) =
676 check_descriptor(store, shards, repo, &context.target(), &key, &raw).await
677 {
678 if error.public_message() != "object blocked" {
679 return Err(error);
680 }
681 denied.insert(*id);
682 }
683 }
684 }
685 }
686 for (id, context) in contexts {
687 for dependency in &context.ids {
688 if self::denied(store, dependency)
689 .await
690 .map_err(|_| unavailable())?
691 {
692 denied.insert(id);
693 }
694 }
695 }
696 Ok(denied)
697}
698
699pub async fn require_pack_clear<S: NamespaceStore>(
700 store: &S,
701 shards: &dyn ShardMap,
702 repo: &RepoId,
703 pack: &Hash,
704) -> Result<(), ServerError> {
705 require_pack_clear_with_concurrency(store, shards, repo, pack, 1).await
706}
707
708pub async fn require_pack_clear_for_scanner<S: NamespaceStore>(
714 store: &S,
715 shards: &dyn ShardMap,
716 repo: &RepoId,
717 pack: &Hash,
718) -> Result<(), ServerError> {
719 require_pack_clear_with_concurrency(store, shards, repo, pack, SCANNER_PROOF_CONCURRENCY).await
720}
721
722async fn require_pack_clear_with_concurrency<S: NamespaceStore>(
723 store: &S,
724 shards: &dyn ShardMap,
725 repo: &RepoId,
726 pack: &Hash,
727 concurrency: usize,
728) -> Result<(), ServerError> {
729 let budget = SliceBudget::new(9000);
730 let remote = Budgeted::new(store, &budget);
731 require_clear(&remote, pack).await?;
732 super::inventory::visit(&remote, pack, false, |_, _| async { Ok(false) })
733 .await
734 .map_err(|_| unavailable())?;
735 let start = Key::new(INDEX_PREFIX.to_vec());
736 let mut end = INDEX_PREFIX.to_vec();
737 *end.last_mut().ok_or_else(unavailable)? = 1;
738 let end = Key::new(end);
739 let target = Target {
740 ids: &BTreeSet::from([*pack]),
741 manifests: &[],
742 packs: vec![*pack],
743 };
744 prove_shards(&remote, shards, repo, &target, &start, &end, concurrency).await?;
745 for id in target.ids {
746 require_clear(&remote, id).await?;
747 }
748 Ok(())
749}
750
751#[cfg(test)]
752mod tests {
753 use super::*;
754 use crate::pipeline::D34Shards;
755 use crate::store::{
756 BatchOutcome, Cursor, Partition, PartitionStats, ScanPage, StoreCapabilities,
757 };
758 use crate::store::{
759 content_shard,
760 index::{IndexEntry, IndexValue},
761 };
762 use crate::{Batch, ManualClock, MemoryKv, NamespaceKey, RepoName};
763 use futures_executor::block_on;
764 use std::sync::Arc;
765 use std::sync::atomic::{AtomicUsize, Ordering};
766
767 struct ScanProbe {
768 inner: MemoryKv,
769 active: AtomicUsize,
770 peak: AtomicUsize,
771 scans: AtomicUsize,
772 completed: AtomicUsize,
773 nested: AtomicUsize,
774 fail_prefix: Option<u16>,
775 stall: bool,
776 }
777 impl ScanProbe {
778 fn new(fail_prefix: Option<u16>, stall: bool) -> Self {
779 Self {
780 inner: store(),
781 active: AtomicUsize::new(0),
782 peak: AtomicUsize::new(0),
783 scans: AtomicUsize::new(0),
784 completed: AtomicUsize::new(0),
785 nested: AtomicUsize::new(0),
786 fail_prefix,
787 stall,
788 }
789 }
790 }
791 struct ActiveScan<'a>(&'a AtomicUsize);
792 impl Drop for ActiveScan<'_> {
793 fn drop(&mut self) {
794 self.0.fetch_sub(1, Ordering::SeqCst);
795 }
796 }
797 impl NamespaceStore for ScanProbe {
798 fn capabilities(&self) -> StoreCapabilities {
799 self.inner.capabilities()
800 }
801 async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
802 assert_eq!(
803 self.active.load(Ordering::SeqCst),
804 0,
805 "nested reads await sibling EOF"
806 );
807 self.nested.fetch_add(1, Ordering::SeqCst);
808 self.inner.get(p, key).await
809 }
810 async fn scan(
811 &self,
812 p: &Partition,
813 start: &Key,
814 end: &Key,
815 after: Option<&Cursor>,
816 limit: u32,
817 ) -> Result<ScanPage, StoreError> {
818 if start.as_bytes() != super::super::directory::PREFIX {
819 assert_eq!(
820 self.active.load(Ordering::SeqCst),
821 0,
822 "nested scan awaits sibling EOF"
823 );
824 self.nested.fetch_add(1, Ordering::SeqCst);
825 return self.inner.scan(p, start, end, after, limit).await;
826 }
827 assert_eq!(limit, super::super::directory::PAGE_ROWS);
828 if after.is_some() {
829 assert_eq!(
830 self.active.load(Ordering::SeqCst),
831 0,
832 "continuations await sibling EOF"
833 );
834 }
835 self.scans.fetch_add(1, Ordering::SeqCst);
836 let active = self.active.fetch_add(1, Ordering::SeqCst) + 1;
837 self.peak.fetch_max(active, Ordering::SeqCst);
838 let _guard = ActiveScan(&self.active);
839 let mut remaining = if matches!(p, Partition::ContentShard(prefix) if Some(*prefix) == self.fail_prefix)
840 {
841 1
842 } else {
843 3
844 };
845 futures::future::poll_fn(|cx| {
846 if self.stall {
847 return std::task::Poll::Pending;
848 }
849 if remaining == 0 {
850 std::task::Poll::Ready(())
851 } else {
852 remaining -= 1;
853 cx.waker().wake_by_ref();
854 std::task::Poll::Pending
855 }
856 })
857 .await;
858 self.completed.fetch_add(1, Ordering::SeqCst);
859 if matches!(p, Partition::ContentShard(prefix) if Some(*prefix) == self.fail_prefix) {
860 return Err(StoreError::unavailable("probe shard failure"));
861 }
862 self.inner.scan(p, start, end, after, limit).await
863 }
864 async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
865 self.inner.apply(p, batch).await
866 }
867 async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
868 self.inner.stats(p).await
869 }
870 async fn probe(&self) -> Result<(), StoreError> {
871 self.inner.probe().await
872 }
873 }
874 async fn probe_proof(
875 store: &ScanProbe,
876 budget: &SliceBudget,
877 ids: &BTreeSet<Hash>,
878 concurrency: usize,
879 ) -> Result<(), ServerError> {
880 let mut end = INDEX_PREFIX.to_vec();
881 *end.last_mut().unwrap() = 1;
882 prove_shards(
883 &Budgeted::new(store, budget),
884 &D34Shards,
885 &repo("probe"),
886 &Target {
887 ids,
888 manifests: &[],
889 packs: Vec::new(),
890 },
891 &Key::new(INDEX_PREFIX.to_vec()),
892 &Key::new(end),
893 concurrency,
894 )
895 .await
896 }
897
898 #[test]
899 fn scanner_proof_checks_all_shards_with_measured_concurrency_and_shared_budget() {
900 block_on(async {
901 for concurrency in [1, WRITE_PROOF_CONCURRENCY, SCANNER_PROOF_CONCURRENCY] {
902 let store = ScanProbe::new(None, false);
903 let budget = SliceBudget::new(9000);
904 probe_proof(&store, &budget, &BTreeSet::new(), concurrency)
905 .await
906 .unwrap();
907 assert_eq!(store.peak.load(Ordering::SeqCst), concurrency);
908 assert!(store.peak.load(Ordering::SeqCst) <= 6);
909 assert_eq!(store.scans.load(Ordering::SeqCst), 16);
910 assert_eq!(budget.used(), 16);
911 assert_eq!(store.active.load(Ordering::SeqCst), 0);
912 }
913 });
914 }
915
916 #[test]
917 fn scanner_proof_preserves_late_block_failure_and_continuations() {
918 block_on(async {
919 for concurrency in [1, WRITE_PROOF_CONCURRENCY, SCANNER_PROOF_CONCURRENCY] {
920 let store = ScanProbe::new(None, false);
921 let blocked_id = [255; 32];
922 ContentIndex::new(BorrowedStore(&store.inner))
923 .install_block_action(&blocked_id, &action(1, vec![]), 1)
924 .await
925 .unwrap();
926 assert!(matches!(
927 probe_proof(
928 &store,
929 &SliceBudget::new(9000),
930 &BTreeSet::from([blocked_id]),
931 concurrency
932 )
933 .await,
934 Err(error) if error.code() == crate::Code::PermissionDenied
935 ));
936 assert_eq!(store.scans.load(Ordering::SeqCst), 16);
937 assert_eq!(store.active.load(Ordering::SeqCst), 0);
938
939 let mut second_id = blocked_id;
941 second_id[31] = 254;
942 ContentIndex::new(BorrowedStore(&store.inner))
943 .install_block_action(&second_id, &action(2, vec![]), 1)
944 .await
945 .unwrap();
946 store.scans.store(0, Ordering::SeqCst);
947 probe_proof(
948 &store,
949 &SliceBudget::new(9000),
950 &BTreeSet::new(),
951 concurrency,
952 )
953 .await
954 .unwrap();
955 assert_eq!(store.scans.load(Ordering::SeqCst), 16);
956
957 store
958 .inner
959 .apply(
960 &content_shard(&second_id),
961 Batch::new()
962 .put(descriptor_key(&second_id), Value::new(b"corrupt".to_vec())),
963 )
964 .await
965 .unwrap();
966 store.scans.store(0, Ordering::SeqCst);
967 assert!(matches!(
968 probe_proof(&store, &SliceBudget::new(9000), &BTreeSet::new(), concurrency).await,
969 Err(error) if error.code() == crate::Code::Unavailable
970 ));
971 assert_eq!(store.scans.load(Ordering::SeqCst), 16);
972 assert_eq!(store.active.load(Ordering::SeqCst), 0);
973
974 let failed =
975 ScanProbe::new(Some(super::super::directory::DIRECTORY_SHARDS - 1), false);
976 assert!(matches!(
977 probe_proof(
978 &failed,
979 &SliceBudget::new(9000),
980 &BTreeSet::new(),
981 concurrency
982 )
983 .await,
984 Err(error) if error.code() == crate::Code::Unavailable
985 ));
986 assert_eq!(failed.scans.load(Ordering::SeqCst), 16);
987 assert_eq!(failed.active.load(Ordering::SeqCst), 0);
988 }
989 });
990 }
991
992 #[test]
993 fn proof_nested_chunks_wait_for_sibling_eof_and_preserve_valid_results() {
994 block_on(async {
995 for concurrency in [WRITE_PROOF_CONCURRENCY, SCANNER_PROOF_CONCURRENCY] {
996 let store = ScanProbe::new(None, false);
997 let manifest = [0; 32];
998 let chunk = [8; 32];
999 let r = repo("probe");
1000 ContentIndex::new(BorrowedStore(&store.inner))
1001 .install_block_action(&manifest, &action(1, vec![chunk]), 1)
1002 .await
1003 .unwrap();
1004 member_row(&store.inner, &r, [9; 32], manifest).await;
1005 assert!(matches!(probe_proof(&store, &SliceBudget::new(9000),
1006 &BTreeSet::from([chunk]), concurrency).await,
1007 Err(error) if error.code() == crate::Code::PermissionDenied));
1008 assert_eq!(store.scans.load(Ordering::SeqCst), 16);
1009 assert_eq!(store.completed.load(Ordering::SeqCst), 16);
1010 assert!(store.nested.load(Ordering::SeqCst) > 0);
1011 assert_eq!(store.active.load(Ordering::SeqCst), 0);
1012 store.scans.store(0, Ordering::SeqCst);
1013 let budget = SliceBudget::new(9000);
1014 probe_proof(&store, &budget, &BTreeSet::from([[7; 32]]), concurrency)
1015 .await
1016 .unwrap();
1017 assert_eq!(store.scans.load(Ordering::SeqCst), 16);
1018 assert!(budget.used() > 16 && budget.used() < 9000);
1019 assert_eq!(store.active.load(Ordering::SeqCst), 0);
1020 }
1021 });
1022 }
1023
1024 #[test]
1025 fn proof_error_waits_for_every_dispatched_reply_before_return_or_retry() {
1026 block_on(async {
1027 for concurrency in [WRITE_PROOF_CONCURRENCY, SCANNER_PROOF_CONCURRENCY] {
1028 let store = ScanProbe::new(Some(0), false);
1029 let budget = SliceBudget::new(9000);
1030 for attempt in 1..=2 {
1031 assert!(matches!(
1032 probe_proof(&store, &budget, &BTreeSet::new(), concurrency).await,
1033 Err(error) if error.code() == crate::Code::Unavailable
1034 ));
1035 let expected = attempt * concurrency;
1036 assert_eq!(store.scans.load(Ordering::SeqCst), expected);
1037 assert_eq!(store.completed.load(Ordering::SeqCst), expected);
1038 assert_eq!(store.active.load(Ordering::SeqCst), 0);
1039 assert_eq!(budget.used(), u32::try_from(expected).unwrap());
1040 }
1041 assert_eq!(store.peak.load(Ordering::SeqCst), concurrency);
1042 }
1043 });
1044 }
1045
1046 #[test]
1047 fn scanner_proof_budget_failure_and_cancellation_drop_outstanding_scans() {
1048 block_on(async {
1049 let store = ScanProbe::new(None, false);
1050 let budget = SliceBudget::new(13);
1051 assert!(matches!(
1052 probe_proof(&store, &budget, &BTreeSet::new(), SCANNER_PROOF_CONCURRENCY).await,
1053 Err(error) if error.code() == crate::Code::Unavailable
1054 ));
1055 assert_eq!(budget.used(), 13);
1056 assert_eq!(store.scans.load(Ordering::SeqCst), 13);
1057 assert_eq!(store.active.load(Ordering::SeqCst), 0);
1058
1059 let stalled = ScanProbe::new(None, true);
1060 let budget = SliceBudget::new(9000);
1061 let ids = BTreeSet::new();
1062 let mut proof = Box::pin(probe_proof(
1063 &stalled,
1064 &budget,
1065 &ids,
1066 SCANNER_PROOF_CONCURRENCY,
1067 ));
1068 futures::future::poll_fn(|cx| {
1069 assert!(std::future::Future::poll(proof.as_mut(), cx).is_pending());
1070 std::task::Poll::Ready(())
1071 })
1072 .await;
1073 assert_eq!(
1074 stalled.active.load(Ordering::SeqCst),
1075 SCANNER_PROOF_CONCURRENCY
1076 );
1077 drop(proof);
1078 assert_eq!(stalled.active.load(Ordering::SeqCst), 0);
1079 assert_eq!(
1080 stalled.scans.load(Ordering::SeqCst),
1081 SCANNER_PROOF_CONCURRENCY
1082 );
1083 assert_eq!(
1084 budget.used(),
1085 u32::try_from(SCANNER_PROOF_CONCURRENCY).unwrap()
1086 );
1087 });
1088 }
1089
1090 fn store() -> MemoryKv {
1091 MemoryKv::with_clock(Arc::new(ManualClock::new(0)))
1092 }
1093 fn repo(name: &str) -> RepoId {
1094 RepoId {
1095 namespace: NamespaceKey::deployment_default(),
1096 name: RepoName::new(name).unwrap(),
1097 }
1098 }
1099 fn action(id: u8, chunks: Vec<Hash>) -> BlockAction {
1100 BlockAction {
1101 id: [id; 32],
1102 takedown_id: [id; 32],
1103 reason: "legal".into(),
1104 blocked_at_ms: 1,
1105 chunk_ids: chunks,
1106 }
1107 }
1108 #[test]
1109 fn legacy_producer_and_overlapping_actions_cannot_erase_denial() {
1110 block_on(async {
1111 let store = store();
1112 let index = ContentIndex::new(BorrowedStore(&store));
1113 let object = [7; 32];
1114 index
1115 .install_block_action(&object, &action(1, vec![]), 1)
1116 .await
1117 .unwrap();
1118 index
1119 .install_block_action(&object, &action(2, vec![]), 1)
1120 .await
1121 .unwrap();
1122 index
1123 .block(&object, &BlockEntry::new("legacy", 2), 2)
1124 .await
1125 .unwrap();
1126 index.unblock(&object, 3).await.unwrap();
1127 assert!(index.blocked(&object).await.unwrap().is_some());
1128 let raw = store
1129 .get(&content_shard(&object), &action_key(&object))
1130 .await
1131 .unwrap()
1132 .unwrap();
1133 assert_eq!(decode_actions(Some(&raw)).unwrap().len(), 2);
1134 assert!(
1135 index
1136 .install_block_action(&object, &action(1, vec![]), 4)
1137 .await
1138 .is_ok()
1139 );
1140 let mut changed = action(1, vec![]);
1141 changed.reason = "other".into();
1142 assert!(
1143 index
1144 .install_block_action(&object, &changed, 4)
1145 .await
1146 .is_err()
1147 );
1148 });
1149 }
1150 #[test]
1151 fn directory_proof_calls_are_sixteen_plus_descriptors_and_continuations() {
1152 block_on(async {
1153 for count in [0usize, 7, super::super::directory::PAGE_ROWS as usize + 1] {
1154 let store = store();
1155 let index = ContentIndex::new(BorrowedStore(&store));
1156 for n in 0..count {
1157 let mut id = [0; 32];
1158 id[30..].copy_from_slice(&u16::try_from(n).unwrap().to_be_bytes());
1159 index
1160 .install_block_action(&id, &action(1, vec![]), 1)
1161 .await
1162 .unwrap();
1163 }
1164 let budget = SliceBudget::new(9000);
1165 require_repo_clear(&store, &D34Shards, &repo("a"), &BTreeSet::new(), &budget)
1166 .await
1167 .unwrap();
1168 assert_eq!(
1169 budget.used(),
1170 16 + u32::try_from(count).unwrap()
1171 + u32::from(count > super::super::directory::PAGE_ROWS as usize)
1172 );
1173 }
1174 });
1175 }
1176 #[test]
1177 fn empty_authoritative_scan_has_measured_bound() {
1178 block_on(async {
1179 let store = store();
1180 let budget = SliceBudget::new(9000);
1181 require_repo_clear(
1182 &store,
1183 &D34Shards,
1184 &repo("a"),
1185 &BTreeSet::from([[8; 32]]),
1186 &budget,
1187 )
1188 .await
1189 .unwrap();
1190 assert_eq!(budget.used(), 18);
1193 });
1194 }
1195 #[test]
1196 fn chunk_stop_uses_repository_membership_and_not_global_holders() {
1197 block_on(async {
1198 let store = store();
1199 let manifest = [7; 32];
1200 let chunk = [8; 32];
1201 let pack = [9; 32];
1202 let r = repo("a");
1203 let index = ContentIndex::new(BorrowedStore(&store));
1204 index
1205 .install_block_action(&manifest, &action(1, vec![chunk]), 1)
1206 .await
1207 .unwrap();
1208 let entry = IndexEntry {
1209 object: manifest,
1210 value: IndexValue {
1211 frame_offset: 8,
1212 frame_length: 10,
1213 wire_type: 0,
1214 decoded_size: 20,
1215 chain_depth: 0,
1216 delta_base: None,
1217 },
1218 };
1219 let plan = crate::store::index::plan_index_rows_direct(
1220 &D34Shards,
1221 &r,
1222 &content_shard(&manifest),
1223 &pack,
1224 &[entry],
1225 0,
1226 )
1227 .unwrap();
1228 for batch in plan.direct {
1229 let mut writes = Batch::new();
1230 for (k, v) in batch.puts {
1231 writes = writes.put(k, v);
1232 }
1233 store.apply(&batch.target, writes).await.unwrap();
1234 }
1235 store
1236 .apply(
1237 &D34Shards.membership(&r, &crate::BlobKey::pack(pack)),
1238 Batch::new().put(keys::membership(&r.name, &pack), Value::new(vec![1])),
1239 )
1240 .await
1241 .unwrap();
1242 let result = require_repo_clear(
1243 &store,
1244 &D34Shards,
1245 &r,
1246 &BTreeSet::from([chunk]),
1247 &SliceBudget::new(9000),
1248 )
1249 .await;
1250 assert_eq!(result.unwrap_err().public_message(), "object blocked");
1251 require_repo_clear(
1252 &store,
1253 &D34Shards,
1254 &repo("other"),
1255 &BTreeSet::from([chunk]),
1256 &SliceBudget::new(9000),
1257 )
1258 .await
1259 .unwrap();
1260 assert!(!denied(&store, &chunk).await.unwrap());
1261 });
1262 }
1263 #[test]
1264 fn large_chunk_set_reads_only_selected_verified_page_and_rejects_corruption() {
1265 block_on(async {
1266 let store = store();
1267 let manifest = [7; 32];
1268 let chunks: Vec<Hash> = (0u32..20_000)
1269 .map(|i| {
1270 let mut id = [0; 32];
1271 id[..4].copy_from_slice(&i.to_be_bytes());
1272 id
1273 })
1274 .collect();
1275 let index = ContentIndex::new(BorrowedStore(&store));
1276 index
1277 .install_block_action(&manifest, &action(1, chunks.clone()), 1)
1278 .await
1279 .unwrap();
1280 let raw = store
1281 .get(&content_shard(&manifest), &action_key(&manifest))
1282 .await
1283 .unwrap()
1284 .unwrap();
1285 let rows = decode_actions(Some(&raw)).unwrap();
1286 assert_eq!(rows[0].pages.len(), 2);
1287 let budget = SliceBudget::new(2);
1288 let counted = Budgeted::new(&store, &budget);
1289 assert!(
1290 chunks_intersect(
1291 &counted,
1292 &manifest,
1293 &rows[0],
1294 &BTreeSet::from([chunks[18_000]])
1295 )
1296 .await
1297 .unwrap()
1298 );
1299 assert_eq!(budget.used(), 1);
1300 store
1301 .apply(
1302 &content_shard(&manifest),
1303 Batch::new().put(
1304 chunk_page_key(&manifest, &[1; 32], 1),
1305 Value::new(vec![0; 32]),
1306 ),
1307 )
1308 .await
1309 .unwrap();
1310 assert!(
1311 chunks_intersect(
1312 &store,
1313 &manifest,
1314 &rows[0],
1315 &BTreeSet::from([chunks[18_000]])
1316 )
1317 .await
1318 .is_err()
1319 );
1320 assert!(index.blocked(&manifest).await.unwrap().is_some());
1321 });
1322 }
1323 async fn member_row(store: &MemoryKv, repo: &RepoId, pack: Hash, id: Hash) {
1324 let value = IndexValue {
1325 frame_offset: 1,
1326 frame_length: 1,
1327 wire_type: 0,
1328 decoded_size: 1,
1329 chain_depth: 0,
1330 delta_base: None,
1331 };
1332 let plan = crate::store::index::plan_index_rows_direct(
1333 &D34Shards,
1334 repo,
1335 &D34Shards.ref_shard(repo, "refs/heads/main"),
1336 &pack,
1337 &[IndexEntry { object: id, value }],
1338 1,
1339 )
1340 .unwrap();
1341 for batch in plan.direct {
1342 let mut writes = Batch::new();
1343 for (k, v) in batch.puts {
1344 writes = writes.put(k, v);
1345 }
1346 store.apply(&batch.target, writes).await.unwrap();
1347 }
1348 store
1349 .apply(
1350 &D34Shards.membership(repo, &crate::BlobKey::pack(pack)),
1351 Batch::new().put(keys::membership(&repo.name, &pack), Value::new(vec![1])),
1352 )
1353 .await
1354 .unwrap();
1355 }
1356 #[test]
1357 fn million_chunk_manifest_stages_and_proves_with_bounded_pages() {
1358 use mkit_core::object::{ChunkedBlob, Object};
1359 block_on(async {
1360 let store = store();
1361 let repo = repo("large");
1362 let pack = [90; 32];
1363 let chunks: Vec<Hash> = (0u32..1_000_000)
1364 .rev()
1365 .map(|n| {
1366 let mut id = [0; 32];
1367 id[..4].copy_from_slice(&n.to_be_bytes());
1368 id
1369 })
1370 .collect();
1371 let object = Object::ChunkedBlob(ChunkedBlob {
1372 total_size: 65_536_000_000,
1373 chunk_size: 65_536,
1374 chunks,
1375 });
1376 let wire = mkit_core::serialize::serialize(&object).unwrap();
1378 assert_eq!(mkit_core::serialize::deserialize(&wire).unwrap(), object);
1379 drop(wire);
1380 let id = object.id().unwrap();
1381 super::super::inventory::stage(&store, &pack, 1, &id, &object, None, 1)
1382 .await
1383 .unwrap();
1384 super::super::inventory::complete(&store, &pack, 1, 1)
1385 .await
1386 .unwrap();
1387 let row = super::super::inventory::entry(&store, &pack, &id)
1388 .await
1389 .unwrap()
1390 .unwrap();
1391 assert_eq!(row.references.chunk_count, 1_000_000);
1392 assert_eq!(
1393 row.references.pages.len(),
1394 1_000_000usize.div_ceil(PAGE_HASHES)
1395 );
1396 assert!(!row.references.sorted_pages);
1397 assert!(
1398 STAGED_PAGE_PEAK.load(std::sync::atomic::Ordering::Relaxed)
1399 <= 2 * crate::store::MAX_VALUE_BYTES
1400 );
1401 assert!(row.references.pages.capacity() * std::mem::size_of::<ChunkPage>() < 32 * 1024);
1402 member_row(&store, &repo, pack, id).await;
1403 let budget = SliceBudget::new(9000);
1404 require_object_clear(
1405 &store,
1406 &D34Shards,
1407 &repo,
1408 &id,
1409 &IndexedConfig::default(),
1410 &budget,
1411 )
1412 .await
1413 .unwrap();
1414 assert!(budget.used() < 4200, "{}", budget.used());
1415 let Object::ChunkedBlob(manifest) = &object else {
1417 unreachable!()
1418 };
1419 let shared = *manifest.chunks.last().unwrap();
1420 let blocked_id = [91; 32];
1421 member_row(&store, &repo, pack, blocked_id).await;
1422 ContentIndex::new(BorrowedStore(&store))
1423 .install_block_action(&blocked_id, &action(92, vec![shared]), 1)
1424 .await
1425 .unwrap();
1426 let budget = SliceBudget::new(9000);
1427 let denied = require_object_clear(
1428 &store,
1429 &D34Shards,
1430 &repo,
1431 &id,
1432 &IndexedConfig::default(),
1433 &budget,
1434 )
1435 .await
1436 .unwrap_err();
1437 assert_eq!(denied.public_message(), "object blocked");
1438 assert!(budget.used() < 4200, "{}", budget.used());
1439 });
1440 }
1441 #[test]
1442 fn legacy_hold_release_keeps_chunk_denial_and_unblock_removes_it() {
1443 use mkit_core::object::{ChunkedBlob, Object};
1444 block_on(async {
1445 let store = store();
1446 let repo = repo("legacy-context");
1447 let pack = [81; 32];
1448 let id = [82; 32];
1449 let chunk = [83; 32];
1450 super::super::inventory::stage(
1451 &store,
1452 &pack,
1453 1,
1454 &id,
1455 &Object::ChunkedBlob(ChunkedBlob {
1456 total_size: 1,
1457 chunk_size: 0,
1458 chunks: vec![chunk],
1459 }),
1460 None,
1461 1,
1462 )
1463 .await
1464 .unwrap();
1465 super::super::inventory::complete(&store, &pack, 1, 1)
1466 .await
1467 .unwrap();
1468 member_row(&store, &repo, pack, id).await;
1469 let index = ContentIndex::new(BorrowedStore(&store));
1470 index
1471 .block(&id, &BlockEntry::new("policy", 1), 1)
1472 .await
1473 .unwrap();
1474 index.release_hold(&id, &[84; 32], 1).await.unwrap();
1475 let result = require_repo_clear(
1476 &store,
1477 &D34Shards,
1478 &repo,
1479 &BTreeSet::from([chunk]),
1480 &SliceBudget::new(9000),
1481 )
1482 .await
1483 .unwrap_err();
1484 assert_eq!(result.public_message(), "object blocked");
1485 index.unblock(&id, 1).await.unwrap();
1486 require_repo_clear(
1487 &store,
1488 &D34Shards,
1489 &repo,
1490 &BTreeSet::from([chunk]),
1491 &SliceBudget::new(9000),
1492 )
1493 .await
1494 .unwrap();
1495 });
1496 }
1497 #[test]
1498 fn inventory_row_corruption_cannot_change_manifest_type() {
1499 use mkit_core::object::{ChunkedBlob, Object};
1500 block_on(async {
1501 let store = store();
1502 let repo = repo("corrupt");
1503 let pack = [71; 32];
1504 let id = [72; 32];
1505 super::super::inventory::stage(
1506 &store,
1507 &pack,
1508 1,
1509 &id,
1510 &Object::ChunkedBlob(ChunkedBlob {
1511 total_size: 1,
1512 chunk_size: 0,
1513 chunks: vec![[73; 32]],
1514 }),
1515 None,
1516 1,
1517 )
1518 .await
1519 .unwrap();
1520 super::super::inventory::complete(&store, &pack, 1, 1)
1521 .await
1522 .unwrap();
1523 member_row(&store, &repo, pack, id).await;
1524 let mut row = super::super::inventory::entry(&store, &pack, &id)
1525 .await
1526 .unwrap()
1527 .unwrap();
1528 row.kind = 1;
1529 store
1530 .apply(
1531 &content_shard(&pack),
1532 Batch::new().put(
1533 super::super::inventory::entry_key(&pack, &id),
1534 Value::new(serde_json::to_vec(&row).unwrap()),
1535 ),
1536 )
1537 .await
1538 .unwrap();
1539 assert!(
1540 super::super::inventory::member(&store, &D34Shards, &repo, &id)
1541 .await
1542 .is_err()
1543 );
1544 });
1545 }
1546}