use super::denial::{self, BlockAction, StoredAction};
use crate::pipeline::ShardMap;
use crate::store::{Key, StoreError, Value, content_shard, index, keys};
use crate::{Batch, BatchOutcome, NamespaceStore, Precondition, RepoId};
use mkit_core::{
hash::Hash,
object::Object,
ops::graph::{ClosureMode, children},
};
use serde::{Deserialize, Serialize};
use std::collections::BTreeSet;
pub const SCAN_ROWS: u32 = 8;
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Entry {
pub version: u8,
pub kind: u8,
pub canonical_len: u64,
pub logical_len: Option<u64>,
pub base: Option<Hash>,
pub references: StoredAction,
}
impl Entry {
fn validate(&self) -> Result<(), StoreError> {
if self.version != 1
|| self.kind > 7
|| match self.kind {
0 => self.canonical_len != 0 || self.logical_len.is_some(),
1 => {
self.canonical_len < 10
|| self.logical_len != self.canonical_len.checked_sub(10)
}
5 => self.canonical_len < 22 || self.logical_len.is_none(),
_ => self.canonical_len < 6 || self.logical_len.is_some(),
}
{
return Err(bad());
}
denial::encode_actions(vec![self.references.clone()])?;
Ok(())
}
}
#[derive(Default, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
struct Head {
version: u8,
length: u64,
count: u64,
parents: u64,
parent_digest: Hash,
digest: Hash,
complete: bool,
packlist: Option<PacklistFacts>,
}
#[derive(Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
struct PacklistFacts {
prev: Option<Hash>,
packs: Vec<Hash>,
}
pub async fn stage_packlist<S: NamespaceStore>(
store: &S,
pack: &Hash,
length: u64,
prev: Option<Hash>,
packs: &[Hash],
now: u64,
) -> Result<(), StoreError> {
if packs.len()
> crate::store::index::MAX_LOOKUP_IDS + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
{
return Err(bad());
}
let p = content_shard(pack);
let key = head_key(pack);
let raw = store.get(&p, &key).await?;
let mut head: Head = raw.as_ref().map(decode).transpose()?.unwrap_or(Head {
version: 1,
length,
..Head::default()
});
let facts = PacklistFacts {
prev,
packs: packs.to_vec(),
};
if head.version != 1 || head.length != length {
return Err(bad());
}
if let Some(old) = head.packlist {
return if old == facts { Ok(()) } else { Err(bad()) };
}
if head.complete {
return Err(bad());
}
head.packlist = Some(facts);
if store
.apply(
&p,
Batch::new()
.require(guard(key.clone(), raw))
.require(deadline(now))
.put(key, encode(&head)?),
)
.await?
!= BatchOutcome::Committed
{
return Err(StoreError::unavailable("packlist inventory contention"));
}
Ok(())
}
pub(crate) async fn packlist_facts<S: NamespaceStore>(
store: &S,
pack: &Hash,
) -> Result<(u64, Option<Hash>, Vec<Hash>), StoreError> {
let raw = store
.get(&content_shard(pack), &head_key(pack))
.await?
.ok_or_else(bad)?;
let head: Head = decode(&raw)?;
if head.version != 1 || !head.complete {
return Err(bad());
}
let facts = head.packlist.ok_or_else(bad)?;
if facts.packs.len()
> crate::store::index::MAX_LOOKUP_IDS + crate::store::outbox::MAX_TICKETS_PER_ADVANCE
{
return Err(bad());
}
Ok((head.length, facts.prev, facts.packs))
}
#[must_use]
pub fn entry_key(pack: &Hash, id: &Hash) -> Key {
Key::new([keys::block(pack).as_bytes(), b"\0inventory\0", id].concat())
}
#[must_use]
pub fn marker_key(pack: &Hash, id: &Hash) -> Key {
Key::new([keys::block(pack).as_bytes(), b"\0inventory-seal\0", id].concat())
}
fn head_key(pack: &Hash) -> Key {
Key::new([keys::block(pack).as_bytes(), b"\0inventory-head"].concat())
}
fn parent_key(pack: &Hash, id: &Hash) -> Key {
Key::new([keys::block(pack).as_bytes(), b"\0inventory-parent\0", id].concat())
}
fn bad() -> StoreError {
StoreError::Corrupt("invalid verified pack inventory".into())
}
fn encode<T: Serialize>(v: &T) -> Result<Value, StoreError> {
let raw = serde_json::to_vec(v).map_err(|_| bad())?;
if raw.len() > crate::store::MAX_VALUE_BYTES {
return Err(StoreError::Full);
}
Ok(Value::new(raw))
}
fn decode<T: serde::de::DeserializeOwned>(v: &Value) -> Result<T, StoreError> {
serde_json::from_slice(v.as_bytes()).map_err(|_| bad())
}
fn guard(key: Key, raw: Option<Value>) -> Precondition {
match raw {
None => Precondition::Absent(key),
Some(v) => Precondition::Equals(key, v),
}
}
fn deadline(now: u64) -> Precondition {
Precondition::NotAfter(now.saturating_add(crate::store::CONTENT_APPLY_WINDOW_MS))
}
fn entry_digest(id: &Hash, row: &Value) -> Hash {
let mut hash = blake3::Hasher::new();
hash.update(id);
hash.update(row.as_bytes());
*hash.finalize().as_bytes()
}
fn add_digest(digest: &mut Hash, id: &Hash, row: &Value) {
let hash = entry_digest(id, row);
for (a, b) in digest.iter_mut().zip(hash) {
*a ^= b;
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(super) struct InventoryCursor {
after: Option<Vec<u8>>,
count: u64,
digest: Hash,
}
pub(super) async fn has_seal<S: NamespaceStore>(
store: &S,
pack: &Hash,
) -> Result<bool, StoreError> {
let Some(raw) = store.get(&content_shard(pack), &head_key(pack)).await? else {
return Ok(false);
};
let head: Head = decode(&raw)?;
if head.version != 1 || !head.complete {
return Err(bad());
}
Ok(true)
}
pub(super) async fn next<S: NamespaceStore>(
store: &S,
pack: &Hash,
mut state: InventoryCursor,
) -> Result<(InventoryCursor, Vec<(Hash, Entry)>, bool), StoreError> {
let raw = store
.get(&content_shard(pack), &head_key(pack))
.await?
.ok_or_else(bad)?;
let head: Head = decode(&raw)?;
if head.version != 1 || !head.complete || state.count > head.count {
return Err(bad());
}
let start = Key::new([keys::block(pack).as_bytes(), b"\0inventory\0"].concat());
let mut end = start.as_bytes().to_vec();
*end.last_mut().ok_or_else(bad)? = 1;
let after = state.after.clone().map(crate::Cursor::new);
let page = store
.scan(
&content_shard(pack),
&start,
&Key::new(end),
after.as_ref(),
SCAN_ROWS,
)
.await?;
let mut entries = Vec::new();
for (key, raw) in page.entries {
let id: Hash = key
.as_bytes()
.strip_prefix(start.as_bytes())
.ok_or_else(bad)?
.try_into()
.map_err(|_| bad())?;
if store
.get(&content_shard(pack), &marker_key(pack, &id))
.await?
.as_ref()
.map(Value::as_bytes)
!= Some(entry_digest(&id, &raw).as_slice())
{
return Err(bad());
}
let entry: Entry = decode(&raw)?;
if entry.version != 1 || entry.kind > 7 {
return Err(bad());
}
denial::encode_actions(vec![entry.references.clone()])?;
state.count = state.count.checked_add(1).ok_or_else(bad)?;
add_digest(&mut state.digest, &id, &raw);
entries.push((id, entry));
}
if page
.next
.as_ref()
.is_some_and(|next| after.as_ref() == Some(next))
{
return Err(bad());
}
state.after = page.next.map(|next| next.as_bytes().to_vec());
let done = state.after.is_none();
if state.count > head.count || done && (state.count, state.digest) != (head.count, head.digest)
{
return Err(bad());
}
Ok((state, entries, done))
}
pub async fn stage<S: NamespaceStore>(
store: &S,
pack: &Hash,
length: u64,
id: &Hash,
object: &Object,
base: Option<Hash>,
now: u64,
) -> Result<(), StoreError> {
let p = content_shard(pack);
let key = entry_key(pack, id);
if let Some(raw) = store.get(&p, &key).await? {
let existing: Entry = decode(&raw)?;
if existing.kind != 0 {
return Ok(());
}
}
let history_refs: Vec<Hash> = children(object, ClosureMode::History).into_iter().collect();
let chunks = history_refs.as_slice();
let (canonical_len, logical_len) = match object {
Object::Blob(b) => (
(b.data.len() as u64).checked_add(10).ok_or_else(bad)?,
Some(b.data.len() as u64),
),
Object::ChunkedBlob(cb) => (
(cb.chunks.len() as u64)
.checked_mul(32)
.and_then(|n| n.checked_add(22))
.ok_or_else(bad)?,
Some(cb.total_size),
),
_ => (
mkit_core::serialize::serialize(object)
.map_err(|_| bad())?
.len() as u64,
None,
),
};
let references = denial::stage_references(
store,
pack,
&BlockAction {
id: *id,
takedown_id: *pack,
reason: "inventory".into(),
blocked_at_ms: 0,
chunk_ids: Vec::new(),
},
chunks,
false,
now,
)
.await?;
let row = encode(&Entry {
version: 1,
kind: object.object_type() as u8,
canonical_len,
logical_len,
base,
references,
})?;
put_entry(store, pack, length, id, row, now).await
}
async fn put_entry<S: NamespaceStore>(
store: &S,
pack: &Hash,
length: u64,
id: &Hash,
row: Value,
now: u64,
) -> Result<(), StoreError> {
let p = content_shard(pack);
let key = entry_key(pack, id);
let existing = store.get(&p, &key).await?;
if let Some(old) = &existing {
let prior: Entry = decode(old)?;
let next: Entry = decode(&row)?;
if prior.kind != 0 || next.kind == 0 {
return Ok(());
}
}
let hk = head_key(pack);
let old = store.get(&p, &hk).await?;
let mut head: Head = old.as_ref().map(decode).transpose()?.unwrap_or(Head {
version: 1,
length,
..Head::default()
});
if head.version != 1 || head.length != length || head.complete {
return Err(bad());
}
if let Some(old) = &existing {
add_digest(&mut head.digest, id, old);
} else {
head.count = head.count.checked_add(1).ok_or_else(bad)?;
}
add_digest(&mut head.digest, id, &row);
let parent: Entry = decode(&row)?;
if matches!(parent.kind, 2 | 5) {
head.parents = head.parents.checked_add(1).ok_or_else(bad)?;
add_digest(&mut head.parent_digest, id, &row);
}
let mut batch = Batch::new()
.require(guard(hk.clone(), old))
.require(guard(key.clone(), existing))
.require(deadline(now))
.put(hk, encode(&head)?)
.put(key, row.clone())
.put(
marker_key(pack, id),
Value::new(entry_digest(id, &row).to_vec()),
);
if matches!(parent.kind, 2 | 5) {
batch = batch.put(parent_key(pack, id), row);
}
if store.apply(&p, batch).await? != BatchOutcome::Committed {
return Err(StoreError::unavailable("inventory stage contention"));
}
Ok(())
}
pub async fn dependency<S: NamespaceStore>(
store: &S,
pack: &Hash,
length: u64,
id: &Hash,
now: u64,
) -> Result<(), StoreError> {
let refs = denial::stage_action(
store,
pack,
&BlockAction {
id: *id,
takedown_id: *pack,
reason: "inventory".into(),
blocked_at_ms: 0,
chunk_ids: Vec::new(),
},
now,
)
.await?;
put_entry(
store,
pack,
length,
id,
encode(&Entry {
version: 1,
kind: 0,
canonical_len: 0,
logical_len: None,
base: None,
references: refs,
})?,
now,
)
.await
}
pub async fn complete<S: NamespaceStore>(
store: &S,
pack: &Hash,
length: u64,
now: u64,
) -> Result<(), StoreError> {
let p = content_shard(pack);
let key = head_key(pack);
let old = store.get(&p, &key).await?;
let mut head: Head = old.as_ref().map(decode).transpose()?.unwrap_or(Head {
version: 1,
length,
..Head::default()
});
if head.version != 1 || head.length != length {
return Err(bad());
}
if head.complete {
return Ok(());
}
head.complete = true;
let batch = Batch::new()
.require(guard(key.clone(), old))
.require(deadline(now))
.put(key, encode(&head)?);
if store.apply(&p, batch).await? != BatchOutcome::Committed {
return Err(StoreError::unavailable("inventory seal contention"));
}
Ok(())
}
pub async fn seal<S: NamespaceStore>(store: &S, pack: &Hash) -> Result<Hash, StoreError> {
let raw = store
.get(&content_shard(pack), &head_key(pack))
.await?
.ok_or_else(bad)?;
let head: Head = decode(&raw)?;
if head.version != 1 || !head.complete {
return Err(bad());
}
Ok(mkit_core::hash::hash(raw.as_bytes()))
}
pub async fn entry<S: NamespaceStore>(
store: &S,
pack: &Hash,
id: &Hash,
) -> Result<Option<Entry>, StoreError> {
let values = store
.get_many(
&content_shard(pack),
&[entry_key(pack, id), marker_key(pack, id)],
)
.await?;
if values.len() != 2 {
return Err(bad());
}
let raw = values[0].as_ref();
if let Some(raw) = raw
&& values[1].as_ref().map(Value::as_bytes) != Some(entry_digest(id, raw).as_slice())
{
return Err(bad());
}
let row: Option<Entry> = raw.map(decode).transpose()?;
if let Some(row) = &row {
row.validate()?;
}
Ok(row)
}
pub async fn member<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
id: &Hash,
) -> Result<(Hash, Entry), StoreError> {
let rows = index::locate_many(store, shards, repo, &[*id]).await?;
let loc = rows
.first()
.and_then(|v| v.as_ref().ok())
.and_then(|v| *v)
.ok_or_else(bad)?;
seal(store, &loc.pack).await?;
Ok((
loc.pack,
entry(store, &loc.pack, id).await?.ok_or_else(bad)?,
))
}
pub async fn prepare_object<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
id: &Hash,
action: BlockAction,
_now: u64,
) -> Result<StoredAction, StoreError> {
let (_, row) = member(store, shards, repo, id).await?;
if !matches!(row.kind, 1 | 5) {
return Err(StoreError::Invalid(
"takedown target must be file content".into(),
));
}
let mut stored = row.references;
stored.action = action;
if row.kind != 5 {
stored.pages.clear();
stored.chunk_count = 0;
stored.chunk_digest = mkit_core::hash::hash(&[]);
}
Ok(stored)
}
pub async fn prepare_pack<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
pack: &Hash,
action: BlockAction,
now: u64,
) -> Result<StoredAction, StoreError> {
if !crate::store::read::is_member(store, shards, repo, pack, None).await? {
return Err(bad());
}
let digest = seal(store, pack).await?;
let mut stored = denial::stage_action(store, pack, &action, now).await?;
stored.pack_scope = Some(*pack);
stored.pack_digest = Some(digest);
Ok(stored)
}
pub async fn visit<S: NamespaceStore, F, Fut>(
store: &S,
pack: &Hash,
parents: bool,
mut f: F,
) -> Result<bool, StoreError>
where
F: FnMut(Hash, Entry) -> Fut,
Fut: std::future::Future<Output = Result<bool, StoreError>>,
{
let raw = store
.get(&content_shard(pack), &head_key(pack))
.await?
.ok_or_else(bad)?;
let head: Head = decode(&raw)?;
if head.version != 1 || !head.complete {
return Err(bad());
}
let prefix = if parents {
b"\0inventory-parent\0".as_slice()
} else {
b"\0inventory\0".as_slice()
};
let start = Key::new([keys::block(pack).as_bytes(), prefix].concat());
let mut end = start.as_bytes().to_vec();
*end.last_mut().ok_or_else(bad)? = 1;
let end = Key::new(end);
let mut after = None;
let mut count = 0u64;
let mut digest = [0; 32];
loop {
let page = store
.scan(
&content_shard(pack),
&start,
&end,
after.as_ref(),
SCAN_ROWS,
)
.await?;
let marker_keys: Result<Vec<Key>, StoreError> = page
.entries
.iter()
.map(|(key, _)| {
let id: Hash = key
.as_bytes()
.strip_prefix(start.as_bytes())
.ok_or_else(bad)?
.try_into()
.map_err(|_| bad())?;
Ok(marker_key(pack, &id))
})
.collect();
let markers = store.get_many(&content_shard(pack), &marker_keys?).await?;
if markers.len() != page.entries.len() {
return Err(bad());
}
for ((key, raw), marker) in page.entries.into_iter().zip(markers) {
let id: Hash = key
.as_bytes()
.strip_prefix(start.as_bytes())
.ok_or_else(bad)?
.try_into()
.map_err(|_| bad())?;
if marker.as_ref().map(Value::as_bytes) != Some(entry_digest(&id, &raw).as_slice()) {
return Err(bad());
}
let row: Entry = decode(&raw)?;
row.validate()?;
count += 1;
add_digest(&mut digest, &id, &raw);
if f(id, row).await? {
return Ok(true);
}
}
match page.next {
Some(next) if after.as_ref() != Some(&next) => after = Some(next),
Some(_) => return Err(bad()),
None => break,
}
}
let expected = if parents {
(head.parents, head.parent_digest)
} else {
(head.count, head.digest)
};
if (count, digest) != expected {
return Err(bad());
}
Ok(false)
}
pub async fn is_file<S: NamespaceStore>(
store: &S,
pack: &Hash,
id: &Hash,
) -> Result<bool, StoreError> {
let Some(row) = entry(store, pack, id).await? else {
return Ok(false);
};
if row.kind == 5 {
return Ok(true);
}
if row.kind != 1 {
return Ok(false);
}
let chunks = visit(store, pack, true, |_, parent| async move {
if parent.kind != 5 {
return Ok(false);
}
denial::chunks_intersect(store, pack, &parent.references, &BTreeSet::from([*id]))
.await
.map_err(|_| bad())
})
.await?;
if !chunks {
return Ok(true);
}
visit(store, pack, true, |_, parent| async move {
if parent.kind != 2 {
return Ok(false);
}
denial::chunks_intersect(store, pack, &parent.references, &BTreeSet::from([*id]))
.await
.map_err(|_| bad())
})
.await
}