mkit-server 0.5.0

Runtime-agnostic core of the mkit server: operation model, errors, runtime and telemetry vocabulary
Documentation
//! Privileged canonical acquisition with explicit runtime admission geometry.
use crate::indexed::{
    geometry,
    resolve::{self, MemberCache, MemberSourceLimits},
};
use crate::pipeline::ShardMap;
use crate::{BlobStore, Metrics, NamespaceStore, RepoId, ServerError};
use mkit_core::hash::Hash;
use std::{collections::BTreeSet, sync::Arc};

/// Validated limits; the caller budgets the actual namespace and blob boundaries.
#[derive(Debug, Clone, Copy)]
pub struct Profile {
    pub(super) limits: MemberSourceLimits,
    pub(super) chain_depth: u32,
    retained: u64,
    retain_latest: bool,
    resident: u64,
    slice_calls: u32,
}
impl Profile {
    /// Inline verification's configured retained budget and checked resident bound.
    ///
    /// # Errors
    /// Invalid chain geometry or overflow of the conservative resident allowance.
    pub fn inline(decode_budget: u64, chain_depth: u32) -> Result<Self, ServerError> {
        if decode_budget < 8 || chain_depth == 0 || chain_depth > u32::from(u16::MAX) {
            return Err(ServerError::invalid_argument("invalid acquisition limits"));
        }
        let resident = decode_budget
            .checked_mul(8)
            .and_then(|n| n.checked_add(128 << 20))
            .ok_or_else(|| ServerError::invalid_argument("acquisition resident bound overflow"))?;
        Ok(Self {
            limits: MemberSourceLimits {
                max_frame_bytes: decode_budget.min(geometry::FRAME_BYTES),
                max_decoded_bytes: decode_budget.min(geometry::CANONICAL_BYTES),
            },
            chain_depth,
            retained: decode_budget,
            retain_latest: false,
            resident,
            slice_calls: ((chain_depth + 1) * 8 + 256).max(700),
        })
    }
    /// Current Worker admission: 1 MiB payloads plus canonical framing, 50 delta hops.
    /// Resident bound: 16 MiB encoded payload + 28 MiB R-203 decoder scratch
    /// + two canonical entry buffers leaves almost 2 MiB headroom within 48 MiB.
    /// Decoder scratch drops before base copies, delta output, object parsing
    /// and Arc conversion. That later phase fits eight canonical entry regions
    /// alongside the frame, below 25 MiB including metadata. Range collection
    /// and its transport copies drop before decoding; even two full-frame
    /// copies plus the canonical base fit below 34 MiB. The five-byte header fits headroom.
    /// Source-selection descriptors are checkpointed separately across slices.
    #[must_use]
    pub const fn scheduled() -> Self {
        Self {
            limits: MemberSourceLimits {
                max_frame_bytes: geometry::FRAME_BYTES,
                max_decoded_bytes: geometry::CANONICAL_BYTES,
            },
            chain_depth: 50,
            retained: geometry::CANONICAL_BYTES,
            retain_latest: true,
            resident: geometry::RESIDENT_BYTES,
            slice_calls: 700,
        }
    }
    /// Conservative per-acquisition resident allowance for startup reporting.
    #[must_use]
    pub const fn resident_upper_bound(&self) -> u64 {
        self.resident
    }
    pub(super) const fn slice_calls(&self) -> u32 {
        self.slice_calls
    }
}

/// Canonical storable bytes and their exact kind, verified against membership.
#[derive(Debug)]
pub struct Verified {
    pub id: Hash,
    pub kind: u8,
    pub canonical: Arc<[u8]>,
}

/// Resolve using the internal privilege that preserves already denied content.
/// All I/O must receive the caller's shared budget wrappers; retries do not reset it.
///
/// # Errors
/// Unknown membership, non-storable kinds, corrupt sources or exhausted bounds fail closed.
#[allow(clippy::too_many_arguments)]
pub async fn resolve<B: BlobStore, S: NamespaceStore>(
    blobs: &B,
    store: &S,
    shards: &dyn ShardMap,
    repo: &RepoId,
    id: Hash,
    profile: &Profile,
    metrics: &dyn Metrics,
) -> Result<Verified, ServerError> {
    let located = resolve::locate_split(store, shards, repo, &[id], metrics)
        .await?
        .remove(&id)
        .and_then(Result::ok)
        .flatten()
        .ok_or_else(|| ServerError::unavailable("canonical member source unavailable"))?;
    let mut memo = MemberCache::with_work_budget(profile.chain_depth + 1);
    decode(
        blobs, store, shards, repo, id, located, profile, metrics, &mut memo,
    )
    .await
}

#[allow(clippy::too_many_arguments)]
pub(crate) async fn resolve_selected<B: BlobStore, S: NamespaceStore>(
    blobs: &B,
    store: &S,
    shards: &dyn ShardMap,
    repo: &RepoId,
    id: Hash,
    profile: &Profile,
    root: &crate::Partition,
    prefix: &crate::Key,
) -> Result<Verified, ServerError> {
    let raw = store
        .get(
            root,
            &crate::Key::new([prefix.as_bytes(), &0u32.to_be_bytes()].concat()),
        )
        .await
        .map_err(|_| ServerError::unavailable("selected source unavailable"))?
        .ok_or_else(|| ServerError::unavailable("selected source unavailable"))?;
    let (found, located) = super::source::decode_frame(&raw)
        .map_err(|_| ServerError::unavailable("selected source unavailable"))?;
    if found != id {
        return Err(ServerError::unavailable("selected source unavailable"));
    }
    let mut memo =
        MemberCache::with_selection(profile.chain_depth + 1, root.clone(), prefix.clone());
    decode(
        blobs,
        store,
        shards,
        repo,
        id,
        located,
        profile,
        &crate::NoopMetrics,
        &mut memo,
    )
    .await
}

#[allow(clippy::too_many_arguments)]
async fn decode<B: BlobStore, S: NamespaceStore>(
    blobs: &B,
    store: &S,
    shards: &dyn ShardMap,
    repo: &RepoId,
    id: Hash,
    located: crate::store::index::LocatedObject,
    profile: &Profile,
    metrics: &dyn Metrics,
    memo: &mut MemberCache,
) -> Result<Verified, ServerError> {
    if profile.retain_latest {
        memo.retain_latest();
    }
    let (canonical, _) = resolve::member_object_for_preservation_bounded(
        blobs,
        store,
        shards,
        repo,
        id,
        located,
        profile.chain_depth,
        profile.retained,
        memo,
        &mut BTreeSet::new(),
        metrics,
        profile.limits,
    )
    .await
    .map_err(|failure| match failure {
        resolve::ResolveFailure::Corrupt(_) => {
            ServerError::new(crate::Code::DataLoss, "canonical member source corrupt")
        }
        _ => ServerError::unavailable("canonical member source unavailable"),
    })?;
    let kind = canonical.first().copied().unwrap_or(0);
    if !matches!(kind, 1 | 2 | 3 | 4 | 5 | 7) {
        return Err(ServerError::invalid_argument(
            "preservation target must be a canonical storable object",
        ));
    }
    Ok(Verified {
        id,
        kind,
        canonical,
    })
}

#[cfg(test)]
#[path = "acquisition_tests.rs"]
mod tests;