Skip to main content

mkit_server/takedown/
inventory.rs

1//! Immutable typed pack facts, staged during verification and sealed before `vs`.
2use super::denial::{self, BlockAction, StoredAction};
3use crate::pipeline::ShardMap;
4use crate::store::{Key, StoreError, Value, content_shard, index, keys};
5use crate::{Batch, BatchOutcome, NamespaceStore, Precondition, RepoId};
6use mkit_core::{
7    hash::Hash,
8    object::Object,
9    ops::graph::{ClosureMode, children},
10};
11use serde::{Deserialize, Serialize};
12use std::collections::BTreeSet;
13/// Bound scan values and their temporary base64/JSON transport representation.
14pub const SCAN_ROWS: u32 = 8;
15#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
16#[serde(deny_unknown_fields)]
17pub struct Entry {
18    pub version: u8,
19    pub kind: u8,
20    pub canonical_len: u64,
21    pub logical_len: Option<u64>,
22    pub base: Option<Hash>,
23    pub references: StoredAction,
24}
25impl Entry {
26    fn validate(&self) -> Result<(), StoreError> {
27        if self.version != 1
28            || self.kind > 7
29            || match self.kind {
30                0 => self.canonical_len != 0 || self.logical_len.is_some(),
31                1 => {
32                    self.canonical_len < 10
33                        || self.logical_len != self.canonical_len.checked_sub(10)
34                }
35                5 => self.canonical_len < 22 || self.logical_len.is_none(),
36                _ => self.canonical_len < 6 || self.logical_len.is_some(),
37            }
38        {
39            return Err(bad());
40        }
41        denial::encode_actions(vec![self.references.clone()])?;
42        Ok(())
43    }
44}
45#[derive(Default, Clone, Serialize, Deserialize)]
46#[serde(deny_unknown_fields)]
47struct Head {
48    version: u8,
49    length: u64,
50    count: u64,
51    parents: u64,
52    parent_digest: Hash,
53    digest: Hash,
54    complete: bool,
55    packlist: Option<PacklistFacts>,
56}
57#[derive(Clone, Serialize, Deserialize, PartialEq, Eq)]
58#[serde(deny_unknown_fields)]
59struct PacklistFacts {
60    prev: Option<Hash>,
61    packs: Vec<Hash>,
62}
63/// Bind decoded MKPL facts before the verified inventory seal.
64pub async fn stage_packlist<S: NamespaceStore>(
65    store: &S,
66    pack: &Hash,
67    length: u64,
68    prev: Option<Hash>,
69    packs: &[Hash],
70    now: u64,
71) -> Result<(), StoreError> {
72    if packs.len()
73        > crate::store::index::MAX_LOOKUP_IDS + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
74    {
75        return Err(bad());
76    }
77    let p = content_shard(pack);
78    let key = head_key(pack);
79    let raw = store.get(&p, &key).await?;
80    let mut head: Head = raw.as_ref().map(decode).transpose()?.unwrap_or(Head {
81        version: 1,
82        length,
83        ..Head::default()
84    });
85    let facts = PacklistFacts {
86        prev,
87        packs: packs.to_vec(),
88    };
89    if head.version != 1 || head.length != length {
90        return Err(bad());
91    }
92    if let Some(old) = head.packlist {
93        return if old == facts { Ok(()) } else { Err(bad()) };
94    }
95    if head.complete {
96        return Err(bad());
97    }
98    head.packlist = Some(facts);
99    if store
100        .apply(
101            &p,
102            Batch::new()
103                .require(guard(key.clone(), raw))
104                .require(deadline(now))
105                .put(key, encode(&head)?),
106        )
107        .await?
108        != BatchOutcome::Committed
109    {
110        return Err(StoreError::unavailable("packlist inventory contention"));
111    }
112    Ok(())
113}
114/// Canonical source was verified before these immutable facts were sealed.
115pub(crate) async fn packlist_facts<S: NamespaceStore>(
116    store: &S,
117    pack: &Hash,
118) -> Result<(u64, Option<Hash>, Vec<Hash>), StoreError> {
119    let raw = store
120        .get(&content_shard(pack), &head_key(pack))
121        .await?
122        .ok_or_else(bad)?;
123    let head: Head = decode(&raw)?;
124    if head.version != 1 || !head.complete {
125        return Err(bad());
126    }
127    let facts = head.packlist.ok_or_else(bad)?;
128    if facts.packs.len()
129        > crate::store::index::MAX_LOOKUP_IDS + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
130    {
131        return Err(bad());
132    }
133    Ok((head.length, facts.prev, facts.packs))
134}
135#[must_use]
136pub fn entry_key(pack: &Hash, id: &Hash) -> Key {
137    Key::new([keys::block(pack).as_bytes(), b"\0inventory\0", id].concat())
138}
139#[must_use]
140pub fn marker_key(pack: &Hash, id: &Hash) -> Key {
141    Key::new([keys::block(pack).as_bytes(), b"\0inventory-seal\0", id].concat())
142}
143fn head_key(pack: &Hash) -> Key {
144    Key::new([keys::block(pack).as_bytes(), b"\0inventory-head"].concat())
145}
146fn parent_key(pack: &Hash, id: &Hash) -> Key {
147    Key::new([keys::block(pack).as_bytes(), b"\0inventory-parent\0", id].concat())
148}
149fn bad() -> StoreError {
150    StoreError::Corrupt("invalid verified pack inventory".into())
151}
152fn encode<T: Serialize>(v: &T) -> Result<Value, StoreError> {
153    let raw = serde_json::to_vec(v).map_err(|_| bad())?;
154    if raw.len() > crate::store::MAX_VALUE_BYTES {
155        return Err(StoreError::Full);
156    }
157    Ok(Value::new(raw))
158}
159fn decode<T: serde::de::DeserializeOwned>(v: &Value) -> Result<T, StoreError> {
160    serde_json::from_slice(v.as_bytes()).map_err(|_| bad())
161}
162fn guard(key: Key, raw: Option<Value>) -> Precondition {
163    match raw {
164        None => Precondition::Absent(key),
165        Some(v) => Precondition::Equals(key, v),
166    }
167}
168fn deadline(now: u64) -> Precondition {
169    Precondition::NotAfter(now.saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS))
170}
171fn entry_digest(id: &Hash, row: &Value) -> Hash {
172    let mut hash = blake3::Hasher::new();
173    hash.update(id);
174    hash.update(row.as_bytes());
175    *hash.finalize().as_bytes()
176}
177fn add_digest(digest: &mut Hash, id: &Hash, row: &Value) {
178    let hash = entry_digest(id, row);
179    for (a, b) in digest.iter_mut().zip(hash) {
180        *a ^= b;
181    }
182}
183/// A bounded inventory traversal whose aggregate is checked against the seal.
184#[derive(Debug, Clone, Default, Serialize, Deserialize)]
185#[serde(deny_unknown_fields)]
186pub(super) struct InventoryCursor {
187    after: Option<Vec<u8>>,
188    count: u64,
189    digest: Hash,
190}
191pub(super) async fn has_seal<S: NamespaceStore>(
192    store: &S,
193    pack: &Hash,
194) -> Result<bool, StoreError> {
195    let Some(raw) = store.get(&content_shard(pack), &head_key(pack)).await? else {
196        return Ok(false);
197    };
198    let head: Head = decode(&raw)?;
199    if head.version != 1 || !head.complete {
200        return Err(bad());
201    }
202    Ok(true)
203}
204pub(super) async fn next<S: NamespaceStore>(
205    store: &S,
206    pack: &Hash,
207    mut state: InventoryCursor,
208) -> Result<(InventoryCursor, Vec<(Hash, Entry)>, bool), StoreError> {
209    let raw = store
210        .get(&content_shard(pack), &head_key(pack))
211        .await?
212        .ok_or_else(bad)?;
213    let head: Head = decode(&raw)?;
214    if head.version != 1 || !head.complete || state.count > head.count {
215        return Err(bad());
216    }
217    let start = Key::new([keys::block(pack).as_bytes(), b"\0inventory\0"].concat());
218    let mut end = start.as_bytes().to_vec();
219    *end.last_mut().ok_or_else(bad)? = 1;
220    let after = state.after.clone().map(crate::Cursor::new);
221    let page = store
222        .scan(
223            &content_shard(pack),
224            &start,
225            &Key::new(end),
226            after.as_ref(),
227            SCAN_ROWS,
228        )
229        .await?;
230    let mut entries = Vec::new();
231    for (key, raw) in page.entries {
232        let id: Hash = key
233            .as_bytes()
234            .strip_prefix(start.as_bytes())
235            .ok_or_else(bad)?
236            .try_into()
237            .map_err(|_| bad())?;
238        if store
239            .get(&content_shard(pack), &marker_key(pack, &id))
240            .await?
241            .as_ref()
242            .map(Value::as_bytes)
243            != Some(entry_digest(&id, &raw).as_slice())
244        {
245            return Err(bad());
246        }
247        let entry: Entry = decode(&raw)?;
248        if entry.version != 1 || entry.kind > 7 {
249            return Err(bad());
250        }
251        denial::encode_actions(vec![entry.references.clone()])?;
252        state.count = state.count.checked_add(1).ok_or_else(bad)?;
253        add_digest(&mut state.digest, &id, &raw);
254        entries.push((id, entry));
255    }
256    if page
257        .next
258        .as_ref()
259        .is_some_and(|next| after.as_ref() == Some(next))
260    {
261        return Err(bad());
262    }
263    state.after = page.next.map(|next| next.as_bytes().to_vec());
264    let done = state.after.is_none();
265    if state.count > head.count || done && (state.count, state.digest) != (head.count, head.digest)
266    {
267        return Err(bad());
268    }
269    Ok((state, entries, done))
270}
271/// The first occurrence owns an entry, like the write-once object index.
272/// Pages are staged first; entry, parent descriptor and running seal share one CAS.
273pub async fn stage<S: NamespaceStore>(
274    store: &S,
275    pack: &Hash,
276    length: u64,
277    id: &Hash,
278    object: &Object,
279    base: Option<Hash>,
280    now: u64,
281) -> Result<(), StoreError> {
282    let p = content_shard(pack);
283    let key = entry_key(pack, id);
284    if let Some(raw) = store.get(&p, &key).await? {
285        let existing: Entry = decode(&raw)?;
286        if existing.kind != 0 {
287            return Ok(());
288        }
289    }
290    let history_refs: Vec<Hash> = children(object, ClosureMode::History).into_iter().collect();
291    let chunks = history_refs.as_slice();
292    let (canonical_len, logical_len) = match object {
293        Object::Blob(b) => (
294            (b.data.len() as u64).checked_add(10).ok_or_else(bad)?,
295            Some(b.data.len() as u64),
296        ),
297        Object::ChunkedBlob(cb) => (
298            (cb.chunks.len() as u64)
299                .checked_mul(32)
300                .and_then(|n| n.checked_add(22))
301                .ok_or_else(bad)?,
302            Some(cb.total_size),
303        ),
304        _ => (
305            mkit_core::serialize::serialize(object)
306                .map_err(|_| bad())?
307                .len() as u64,
308            None,
309        ),
310    };
311    let references = denial::stage_references(
312        store,
313        pack,
314        &BlockAction {
315            id: *id,
316            takedown_id: *pack,
317            reason: "inventory".into(),
318            blocked_at_ms: 0,
319            chunk_ids: Vec::new(),
320        },
321        chunks,
322        false,
323        now,
324    )
325    .await?;
326    let row = encode(&Entry {
327        version: 1,
328        kind: object.object_type() as u8,
329        canonical_len,
330        logical_len,
331        base,
332        references,
333    })?;
334    put_entry(store, pack, length, id, row, now).await
335}
336async fn put_entry<S: NamespaceStore>(
337    store: &S,
338    pack: &Hash,
339    length: u64,
340    id: &Hash,
341    row: Value,
342    now: u64,
343) -> Result<(), StoreError> {
344    let p = content_shard(pack);
345    let key = entry_key(pack, id);
346    let existing = store.get(&p, &key).await?;
347    if let Some(old) = &existing {
348        let prior: Entry = decode(old)?;
349        let next: Entry = decode(&row)?;
350        if prior.kind != 0 || next.kind == 0 {
351            return Ok(());
352        }
353    }
354    let hk = head_key(pack);
355    let old = store.get(&p, &hk).await?;
356    let mut head: Head = old.as_ref().map(decode).transpose()?.unwrap_or(Head {
357        version: 1,
358        length,
359        ..Head::default()
360    });
361    if head.version != 1 || head.length != length || head.complete {
362        return Err(bad());
363    }
364    if let Some(old) = &existing {
365        add_digest(&mut head.digest, id, old);
366    } else {
367        head.count = head.count.checked_add(1).ok_or_else(bad)?;
368    }
369    add_digest(&mut head.digest, id, &row);
370    let parent: Entry = decode(&row)?;
371    if matches!(parent.kind, 2 | 5) {
372        head.parents = head.parents.checked_add(1).ok_or_else(bad)?;
373        add_digest(&mut head.parent_digest, id, &row);
374    }
375    let mut batch = Batch::new()
376        .require(guard(hk.clone(), old))
377        .require(guard(key.clone(), existing))
378        .require(deadline(now))
379        .put(hk, encode(&head)?)
380        .put(key, row.clone())
381        .put(
382            marker_key(pack, id),
383            Value::new(entry_digest(id, &row).to_vec()),
384        );
385    if matches!(parent.kind, 2 | 5) {
386        batch = batch.put(parent_key(pack, id), row);
387    }
388    if store.apply(&p, batch).await? != BatchOutcome::Committed {
389        return Err(StoreError::unavailable("inventory stage contention"));
390    }
391    Ok(())
392}
393pub async fn dependency<S: NamespaceStore>(
394    store: &S,
395    pack: &Hash,
396    length: u64,
397    id: &Hash,
398    now: u64,
399) -> Result<(), StoreError> {
400    let refs = denial::stage_action(
401        store,
402        pack,
403        &BlockAction {
404            id: *id,
405            takedown_id: *pack,
406            reason: "inventory".into(),
407            blocked_at_ms: 0,
408            chunk_ids: Vec::new(),
409        },
410        now,
411    )
412    .await?;
413    put_entry(
414        store,
415        pack,
416        length,
417        id,
418        encode(&Entry {
419            version: 1,
420            kind: 0,
421            canonical_len: 0,
422            logical_len: None,
423            base: None,
424            references: refs,
425        })?,
426        now,
427    )
428    .await
429}
430/// Source hash validation, complete decoding, extraction and index delivery precede this seal.
431pub async fn complete<S: NamespaceStore>(
432    store: &S,
433    pack: &Hash,
434    length: u64,
435    now: u64,
436) -> Result<(), StoreError> {
437    let p = content_shard(pack);
438    let key = head_key(pack);
439    let old = store.get(&p, &key).await?;
440    let mut head: Head = old.as_ref().map(decode).transpose()?.unwrap_or(Head {
441        version: 1,
442        length,
443        ..Head::default()
444    });
445    if head.version != 1 || head.length != length {
446        return Err(bad());
447    }
448    if head.complete {
449        return Ok(());
450    }
451    head.complete = true;
452    let batch = Batch::new()
453        .require(guard(key.clone(), old))
454        .require(deadline(now))
455        .put(key, encode(&head)?);
456    if store.apply(&p, batch).await? != BatchOutcome::Committed {
457        return Err(StoreError::unavailable("inventory seal contention"));
458    }
459    Ok(())
460}
461pub async fn seal<S: NamespaceStore>(store: &S, pack: &Hash) -> Result<Hash, StoreError> {
462    let raw = store
463        .get(&content_shard(pack), &head_key(pack))
464        .await?
465        .ok_or_else(bad)?;
466    let head: Head = decode(&raw)?;
467    if head.version != 1 || !head.complete {
468        return Err(bad());
469    }
470    Ok(mkit_core::hash::hash(raw.as_bytes()))
471}
472pub async fn entry<S: NamespaceStore>(
473    store: &S,
474    pack: &Hash,
475    id: &Hash,
476) -> Result<Option<Entry>, StoreError> {
477    let values = store
478        .get_many(
479            &content_shard(pack),
480            &[entry_key(pack, id), marker_key(pack, id)],
481        )
482        .await?;
483    if values.len() != 2 {
484        return Err(bad());
485    }
486    let raw = values[0].as_ref();
487    if let Some(raw) = raw
488        && values[1].as_ref().map(Value::as_bytes) != Some(entry_digest(id, raw).as_slice())
489    {
490        return Err(bad());
491    }
492    let row: Option<Entry> = raw.map(decode).transpose()?;
493    if let Some(row) = &row {
494        row.validate()?;
495    }
496    Ok(row)
497}
498pub async fn member<S: NamespaceStore>(
499    store: &S,
500    shards: &dyn ShardMap,
501    repo: &RepoId,
502    id: &Hash,
503) -> Result<(Hash, Entry), StoreError> {
504    let rows = index::locate_many(store, shards, repo, &[*id]).await?;
505    let loc = rows
506        .first()
507        .and_then(|v| v.as_ref().ok())
508        .and_then(|v| *v)
509        .ok_or_else(bad)?;
510    seal(store, &loc.pack).await?;
511    Ok((
512        loc.pack,
513        entry(store, &loc.pack, id).await?.ok_or_else(bad)?,
514    ))
515}
516pub async fn prepare_object<S: NamespaceStore>(
517    store: &S,
518    shards: &dyn ShardMap,
519    repo: &RepoId,
520    id: &Hash,
521    action: BlockAction,
522    _now: u64,
523) -> Result<StoredAction, StoreError> {
524    let (_, row) = member(store, shards, repo, id).await?;
525    if !matches!(row.kind, 1 | 5) {
526        return Err(StoreError::Invalid(
527            "takedown target must be file content".into(),
528        ));
529    }
530    let mut stored = row.references;
531    stored.action = action;
532    if row.kind != 5 {
533        stored.pages.clear();
534        stored.chunk_count = 0;
535        stored.chunk_digest = mkit_core::hash::hash(&[]);
536    }
537    Ok(stored)
538}
539pub async fn prepare_pack<S: NamespaceStore>(
540    store: &S,
541    shards: &dyn ShardMap,
542    repo: &RepoId,
543    pack: &Hash,
544    action: BlockAction,
545    now: u64,
546) -> Result<StoredAction, StoreError> {
547    if !crate::store::read::is_member(store, shards, repo, pack, None).await? {
548        return Err(bad());
549    }
550    let digest = seal(store, pack).await?;
551    let mut stored = denial::stage_action(store, pack, &action, now).await?;
552    stored.pack_scope = Some(*pack);
553    stored.pack_digest = Some(digest);
554    Ok(stored)
555}
556/// Stream rows with a bounded page, verifying the complete seal before success.
557/// Parent-only traversal is for reference lookup; the complete traversal validates count and digest.
558pub async fn visit<S: NamespaceStore, F, Fut>(
559    store: &S,
560    pack: &Hash,
561    parents: bool,
562    mut f: F,
563) -> Result<bool, StoreError>
564where
565    F: FnMut(Hash, Entry) -> Fut,
566    Fut: std::future::Future<Output = Result<bool, StoreError>>,
567{
568    let raw = store
569        .get(&content_shard(pack), &head_key(pack))
570        .await?
571        .ok_or_else(bad)?;
572    let head: Head = decode(&raw)?;
573    if head.version != 1 || !head.complete {
574        return Err(bad());
575    }
576    let prefix = if parents {
577        b"\0inventory-parent\0".as_slice()
578    } else {
579        b"\0inventory\0".as_slice()
580    };
581    let start = Key::new([keys::block(pack).as_bytes(), prefix].concat());
582    let mut end = start.as_bytes().to_vec();
583    *end.last_mut().ok_or_else(bad)? = 1;
584    let end = Key::new(end);
585    let mut after = None;
586    let mut count = 0u64;
587    let mut digest = [0; 32];
588    loop {
589        let page = store
590            .scan(
591                &content_shard(pack),
592                &start,
593                &end,
594                after.as_ref(),
595                SCAN_ROWS,
596            )
597            .await?;
598        let marker_keys: Result<Vec<Key>, StoreError> = page
599            .entries
600            .iter()
601            .map(|(key, _)| {
602                let id: Hash = key
603                    .as_bytes()
604                    .strip_prefix(start.as_bytes())
605                    .ok_or_else(bad)?
606                    .try_into()
607                    .map_err(|_| bad())?;
608                Ok(marker_key(pack, &id))
609            })
610            .collect();
611        let markers = store.get_many(&content_shard(pack), &marker_keys?).await?;
612        if markers.len() != page.entries.len() {
613            return Err(bad());
614        }
615        for ((key, raw), marker) in page.entries.into_iter().zip(markers) {
616            let id: Hash = key
617                .as_bytes()
618                .strip_prefix(start.as_bytes())
619                .ok_or_else(bad)?
620                .try_into()
621                .map_err(|_| bad())?;
622            if marker.as_ref().map(Value::as_bytes) != Some(entry_digest(&id, &raw).as_slice()) {
623                return Err(bad());
624            }
625            let row: Entry = decode(&raw)?;
626            row.validate()?;
627            count += 1;
628            add_digest(&mut digest, &id, &raw);
629            if f(id, row).await? {
630                return Ok(true);
631            }
632        }
633        match page.next {
634            Some(next) if after.as_ref() != Some(&next) => after = Some(next),
635            Some(_) => return Err(bad()),
636            None => break,
637        }
638    }
639    let expected = if parents {
640        (head.parents, head.parent_digest)
641    } else {
642        (head.count, head.digest)
643    };
644    if (count, digest) != expected {
645        return Err(bad());
646    }
647    Ok(false)
648}
649/// A Blob used solely as a chunk is not file content in a whole-pack action.
650pub async fn is_file<S: NamespaceStore>(
651    store: &S,
652    pack: &Hash,
653    id: &Hash,
654) -> Result<bool, StoreError> {
655    let Some(row) = entry(store, pack, id).await? else {
656        return Ok(false);
657    };
658    if row.kind == 5 {
659        return Ok(true);
660    }
661    if row.kind != 1 {
662        return Ok(false);
663    }
664    let chunks = visit(store, pack, true, |_, parent| async move {
665        if parent.kind != 5 {
666            return Ok(false);
667        }
668        denial::chunks_intersect(store, pack, &parent.references, &BTreeSet::from([*id]))
669            .await
670            .map_err(|_| bad())
671    })
672    .await?;
673    if !chunks {
674        return Ok(true);
675    }
676    visit(store, pack, true, |_, parent| async move {
677        if parent.kind != 2 {
678            return Ok(false);
679        }
680        denial::chunks_intersect(store, pack, &parent.references, &BTreeSet::from([*id]))
681            .await
682            .map_err(|_| bad())
683    })
684    .await
685}