use crate::pipeline::ShardMap;
use crate::repo::RepoId;
use crate::store::{
codec,
index::{self, IndexValue, LocatedObject, LookupError, ObjectLookup},
keys,
};
use crate::telemetry::{METRIC_INDEX_LOOKUP_CAPPED, Metrics};
use crate::{BlobBody, BlobKey, BlobStore, BoxFuture, ByteRange, NamespaceStore, ServerError};
use futures::StreamExt as _;
use mkit_core::hash::Hash;
use mkit_core::pack::{DeltaBaseSource, PackError, decode_frame_with, peek_delta_header};
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
fn unavailable() -> ServerError {
ServerError::unavailable("object storage request failed")
}
pub(crate) const DECODE_BUDGET_MESSAGE: &str = "pack exceeds indexed decode budget";
fn budget_exceeded() -> ServerError {
ServerError::invalid_argument(DECODE_BUDGET_MESSAGE)
}
#[derive(Debug)]
pub enum ResolveFailure {
Missing,
Capped,
Corrupt(ServerError),
Other(ServerError),
}
impl From<ServerError> for ResolveFailure {
fn from(error: ServerError) -> Self {
Self::Other(error)
}
}
impl ResolveFailure {
#[must_use]
pub fn public_error(self, now: u64, created: u64, bound: u64) -> ServerError {
match self {
Self::Missing => missing_base(now, created, bound),
Self::Capped => {
ServerError::failed_precondition("delta base not available in this repository")
}
Self::Other(error) | Self::Corrupt(error) => error,
}
}
}
#[must_use]
pub fn lagged(now_ms: u64, created_at_ms: u64, bound_ms: u64) -> bool {
now_ms.saturating_sub(created_at_ms) < bound_ms
}
#[must_use]
pub fn missing_base(now_ms: u64, created_at_ms: u64, bound_ms: u64) -> ServerError {
if lagged(now_ms, created_at_ms, bound_ms) {
ServerError::unavailable("repository membership not yet visible")
} else {
ServerError::failed_precondition("delta base not available in this repository")
}
}
fn cap_reason(cap: LookupError) -> &'static str {
match cap {
LookupError::TooManyRows => "rows",
LookupError::TooManyPages => "pages",
LookupError::TooManyMembershipReads => "membership_reads",
}
}
pub async fn locate_split<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
ids: &[Hash],
metrics: &dyn Metrics,
) -> Result<BTreeMap<Hash, ObjectLookup>, ServerError> {
locate_split_inner(store, shards, repo, ids, Some(metrics)).await
}
pub async fn locate_split_quiet<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
ids: &[Hash],
) -> Result<BTreeMap<Hash, ObjectLookup>, ServerError> {
locate_split_inner(store, shards, repo, ids, None).await
}
async fn locate_split_inner<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
ids: &[Hash],
metrics: Option<&dyn Metrics>,
) -> Result<BTreeMap<Hash, ObjectLookup>, ServerError> {
let mut todo: Vec<Vec<Hash>> = ids
.chunks(index::MAX_LOOKUP_IDS)
.map(<[Hash]>::to_vec)
.collect();
let mut found = BTreeMap::new();
while let Some(chunk) = todo.pop() {
let answers = index::locate_many(store, shards, repo, &chunk)
.await
.map_err(|_| unavailable())?;
let split = chunk.len() > 1
&& answers.iter().any(|answer| {
matches!(
answer,
Err(LookupError::TooManyPages | LookupError::TooManyMembershipReads)
)
});
if split {
let mid = chunk.len() / 2;
todo.push(chunk[mid..].to_vec());
todo.push(chunk[..mid].to_vec());
continue;
}
for (id, answer) in chunk.into_iter().zip(answers) {
if let (Err(cap), Some(metrics)) = (answer, metrics) {
tracing::error!(reason = cap_reason(cap), "object index lookup capped");
metrics.incr(
METRIC_INDEX_LOOKUP_CAPPED,
&[("reason", cap_reason(cap))],
1,
);
}
found.insert(id, answer);
}
}
Ok(found)
}
pub(super) async fn frame_bytes<B: BlobStore>(
blobs: &B,
pack: Hash,
offset: u64,
length: u64,
budget: u64,
) -> Result<Vec<u8>, ServerError> {
if length == 0 {
return Err(ServerError::invalid_argument("object hash mismatch"));
}
if length > budget {
return Err(budget_exceeded());
}
let end = offset.checked_add(length - 1).ok_or_else(unavailable)?;
let body = blobs
.get(
&BlobKey::pack(pack),
Some(ByteRange {
start: offset,
end_inclusive: end,
}),
)
.await
.map_err(|_| unavailable())?
.ok_or_else(unavailable)?;
let mut bytes = Vec::new();
bytes
.try_reserve_exact(usize::try_from(length).map_err(|_| budget_exceeded())?)
.map_err(|_| budget_exceeded())?;
match body {
BlobBody::Bytes(value) => {
if value.len() as u64 != length {
return Err(unavailable());
}
bytes.extend_from_slice(&value);
}
BlobBody::Stream { mut stream, .. } => {
while let Some(chunk) = stream.next().await {
let chunk = chunk.map_err(|_| unavailable())?;
if (bytes.len() as u64).saturating_add(chunk.len() as u64) > length {
return Err(unavailable());
}
bytes.extend_from_slice(&chunk);
}
}
}
if bytes.len() as u64 != length {
return Err(unavailable());
}
Ok(bytes)
}
struct CachedBase(Option<(Hash, Arc<[u8]>)>);
impl DeltaBaseSource for CachedBase {
const VERIFIED: bool = false;
fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
Ok(self
.0
.as_ref()
.filter(|(base, _)| base == id)
.map(|(_, bytes)| bytes.to_vec()))
}
}
type Location = (Hash, Hash, u64);
pub type ResolvedMember = (Arc<[u8]>, u32);
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct MemberCache {
no_reads: BTreeSet<Hash>,
rows: BTreeMap<Location, ResolvedMember>,
retained_bytes: u64,
retain_latest: bool,
remaining_work: Option<u32>,
selection: Option<(crate::Partition, crate::Key)>,
}
impl MemberCache {
#[cfg(feature = "http-objects")]
pub(crate) fn forbid_reads(&mut self, ids: &BTreeSet<Hash>) {
self.no_reads.clone_from(ids);
}
pub(crate) fn with_selection(limit: u32, root: crate::Partition, prefix: crate::Key) -> Self {
Self {
selection: Some((root, prefix)),
..Self::with_work_budget(limit)
}
}
pub(crate) fn with_work_budget(limit: u32) -> Self {
Self {
remaining_work: Some(limit),
..Self::default()
}
}
pub(crate) fn retain_latest(&mut self) {
self.retain_latest = true;
}
fn available(&self, budget: u64) -> Result<u64, ServerError> {
if self.retain_latest {
Ok(budget)
} else {
budget
.checked_sub(self.retained_bytes)
.ok_or_else(budget_exceeded)
}
}
pub(crate) fn charge_work(&mut self, amount: u32) -> Result<(), ResolveFailure> {
if let Some(remaining) = &mut self.remaining_work {
*remaining = remaining
.checked_sub(amount)
.ok_or(ResolveFailure::Capped)?;
}
Ok(())
}
#[must_use]
pub fn len(&self) -> usize {
self.rows.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.rows.is_empty()
}
#[must_use]
pub fn retained_bytes(&self) -> u64 {
self.retained_bytes
}
pub(crate) fn rows(&self) -> impl Iterator<Item = (&Location, (&Arc<[u8]>, &u32))> {
self.rows
.iter()
.map(|(location, (bytes, depth))| (location, (bytes, depth)))
}
fn insert(
&mut self,
location: Location,
value: ResolvedMember,
budget: u64,
) -> Result<(), ResolveFailure> {
if self.retain_latest {
self.rows.clear();
self.retained_bytes = 0;
}
let used = self
.retained_bytes
.checked_add(value.0.len() as u64)
.ok_or_else(budget_exceeded)?;
if used > budget {
return Err(budget_exceeded().into());
}
self.rows.insert(location, value);
self.retained_bytes = used;
Ok(())
}
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
pub fn member_object<'a, B: BlobStore, S: NamespaceStore>(
blobs: &'a B,
store: &'a S,
shards: &'a dyn ShardMap,
repo: &'a RepoId,
id: Hash,
located: LocatedObject,
cap: u32,
budget: u64,
memo: &'a mut MemberCache,
visiting: &'a mut BTreeSet<Location>,
metrics: &'a dyn Metrics,
) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
member_object_inner(
blobs, store, shards, repo, id, located, cap, budget, memo, visiting, metrics, true, None,
)
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
pub fn member_object_for_preservation<'a, B: BlobStore, S: NamespaceStore>(
blobs: &'a B,
store: &'a S,
shards: &'a dyn ShardMap,
repo: &'a RepoId,
id: Hash,
located: LocatedObject,
cap: u32,
budget: u64,
memo: &'a mut MemberCache,
visiting: &'a mut BTreeSet<Location>,
metrics: &'a dyn Metrics,
) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
member_object_inner(
blobs, store, shards, repo, id, located, cap, budget, memo, visiting, metrics, false, None,
)
}
#[derive(Debug, Clone, Copy)]
pub struct MemberSourceLimits {
pub max_frame_bytes: u64,
pub max_decoded_bytes: u64,
}
#[allow(clippy::too_many_arguments)]
pub fn member_object_for_preservation_bounded<'a, B: BlobStore, S: NamespaceStore>(
blobs: &'a B,
store: &'a S,
shards: &'a dyn ShardMap,
repo: &'a RepoId,
id: Hash,
located: LocatedObject,
cap: u32,
budget: u64,
memo: &'a mut MemberCache,
visiting: &'a mut BTreeSet<Location>,
metrics: &'a dyn Metrics,
limits: MemberSourceLimits,
) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
member_object_inner(
blobs,
store,
shards,
repo,
id,
located,
cap,
budget,
memo,
visiting,
metrics,
false,
Some(limits),
)
}
async fn selected_frame<S: NamespaceStore>(
store: &S,
selection: Option<&(crate::Partition, crate::Key)>,
level: usize,
id: Hash,
) -> Result<Option<LocatedObject>, ServerError> {
let Some((root, prefix)) = selection else {
return Ok(None);
};
let key = crate::Key::new(
[
prefix.as_bytes(),
&u32::try_from(level)
.map_err(|_| unavailable())?
.to_be_bytes(),
]
.concat(),
);
let raw = store
.get(root, &key)
.await
.map_err(|_| unavailable())?
.ok_or_else(unavailable)?;
let (found, selected) =
crate::takedown::source::decode_frame(&raw).map_err(|_| unavailable())?;
if found != id {
return Err(unavailable());
}
Ok(Some(selected))
}
async fn member_base<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
base: Hash,
located: LocatedObject,
metrics: &dyn Metrics,
selected: Option<LocatedObject>,
) -> Result<LocatedObject, ResolveFailure> {
if let Some(selected) = selected {
return Ok(selected);
}
let partition = shards.object_index(repo, &base);
let key = keys::object_index(&repo.name, &base, &located.pack);
let same = store
.get_many(&partition, &[key])
.await
.map_err(|_| unavailable())?;
if same.len() != 1 {
return Err(unavailable().into());
}
let same = same.into_iter().next().flatten();
let next = if let Some(value) = same {
let value = codec::decode_object_index(&base, &value).map_err(|_| unavailable())?;
(value.frame_offset < located.value.frame_offset).then_some(LocatedObject {
pack: located.pack,
value,
})
} else {
None
};
Ok(match next {
Some(next) => next,
None => match locate_split(store, shards, repo, &[base], metrics)
.await?
.remove(&base)
{
Some(Ok(Some(next))) => next,
Some(Err(_)) => return Err(ResolveFailure::Capped),
_ => return Err(ResolveFailure::Missing),
},
})
}
#[cfg(feature = "http-objects")]
pub(crate) async fn member_dependencies_clear<S: NamespaceStore>(
store: &S,
shards: &dyn ShardMap,
repo: &RepoId,
mut id: Hash,
mut located: LocatedObject,
cap: u32,
metrics: &dyn Metrics,
) -> Result<bool, ServerError> {
let mut visiting = BTreeSet::new();
loop {
if crate::takedown::denial::denied(store, &id).await?
|| crate::takedown::denial::denied(store, &located.pack).await?
|| !visiting.insert((id, located.pack, located.value.frame_offset))
{
return Ok(false);
}
let Some(base) = located.value.delta_base else {
return Ok(true);
};
if visiting.len() > usize::try_from(cap).unwrap_or(usize::MAX) {
return Ok(false);
}
located = match member_base(store, shards, repo, base, located, metrics, None).await {
Ok(next) => next,
Err(ResolveFailure::Missing | ResolveFailure::Capped) => return Ok(false),
Err(ResolveFailure::Other(error) | ResolveFailure::Corrupt(error)) => {
return Err(error);
}
};
id = base;
}
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
fn member_object_inner<'a, B: BlobStore, S: NamespaceStore>(
blobs: &'a B,
store: &'a S,
shards: &'a dyn ShardMap,
repo: &'a RepoId,
id: Hash,
located: LocatedObject,
cap: u32,
budget: u64,
memo: &'a mut MemberCache,
visiting: &'a mut BTreeSet<Location>,
metrics: &'a dyn Metrics,
enforce_denial: bool,
source_limits: Option<MemberSourceLimits>,
) -> BoxFuture<'a, Result<ResolvedMember, ResolveFailure>> {
Box::pin(async move {
if memo.no_reads.contains(&id) {
return Err(budget_exceeded().into());
}
if enforce_denial {
crate::takedown::denial::require_clear(store, &id).await?;
crate::takedown::denial::require_clear(store, &located.pack).await?;
}
let source_limits = Some(source_limits.unwrap_or(MemberSourceLimits {
max_frame_bytes: super::geometry::FRAME_BYTES,
max_decoded_bytes: super::geometry::CANONICAL_BYTES,
}));
if source_limits.is_some_and(|limits| {
located.value.frame_length > limits.max_frame_bytes
|| located.value.decoded_size > limits.max_decoded_bytes
}) {
return Err(budget_exceeded().into());
}
let location = (id, located.pack, located.value.frame_offset);
if let Some(selected) =
selected_frame(store, memo.selection.as_ref(), visiting.len(), id).await?
{
if selected != located {
return Err(unavailable().into());
}
if !store
.has(
&shards.membership(repo, &BlobKey::pack(located.pack)),
&keys::membership(&repo.name, &located.pack),
)
.await
.map_err(|_| unavailable())?
{
return Err(ResolveFailure::Missing);
}
}
let available = memo.available(budget)?;
if let Some(value) = memo.rows.get(&location) {
if value.1 > cap {
return Err(ServerError::invalid_argument("delta chain too deep").into());
}
return Ok(value.clone());
}
if !visiting.is_empty() {
memo.charge_work(1)?;
}
if !visiting.insert(location) {
return Err(ServerError::invalid_argument("delta chain too deep").into());
}
let result: Result<ResolvedMember, ResolveFailure> = async {
let IndexValue {
frame_offset,
frame_length,
delta_base,
..
} = located.value;
let prefix = frame_bytes(blobs, located.pack, 0, 8, available).await?;
let version = u32::from_le_bytes(prefix[4..8].try_into().map_err(|_| unavailable())?);
let mut depth = 0;
let mut base_bytes = None;
if let Some(base) = delta_base {
if visiting.len() > usize::try_from(cap).unwrap_or(usize::MAX) {
return Err(ServerError::invalid_argument("delta chain too deep").into());
}
let selected =
selected_frame(store, memo.selection.as_ref(), visiting.len(), base).await?;
let next =
member_base(store, shards, repo, base, located, metrics, selected).await?;
let (canonical, base_depth) = member_object_inner(
blobs,
store,
shards,
repo,
base,
next,
cap,
budget,
memo,
visiting,
metrics,
enforce_denial,
source_limits,
)
.await?;
base_bytes = Some((base, canonical));
depth = base_depth.saturating_add(1);
if depth > cap {
return Err(ServerError::invalid_argument("delta chain too deep").into());
}
}
let available = memo.available(budget)?;
let frame = frame_bytes(
blobs,
located.pack,
frame_offset,
frame_length,
source_limits.map_or(available, |limits| limits.max_frame_bytes),
)
.await?;
if source_limits.is_some() {
let claim = match frame.first() {
Some(0x00) => Some(frame.len().saturating_sub(5) as u64),
Some(0x03) => frame.get(5..9).and_then(|bytes| {
bytes.try_into().ok().map(u32::from_le_bytes).map(u64::from)
}),
Some(0x02) => frame.get(42..46).and_then(|bytes| {
bytes.try_into().ok().map(u32::from_le_bytes).map(u64::from)
}),
Some(0x04) => frame
.get(41..)
.and_then(|bytes| peek_delta_header(bytes).ok())
.map(|(_, result)| u64::from(result)),
_ => None,
};
if frame.first().copied() != Some(located.value.wire_type)
|| claim.is_some_and(|size| size != located.value.decoded_size)
{
return Err(ResolveFailure::Corrupt(ServerError::invalid_argument(
"verified source frame metadata mismatch",
)));
}
}
let mut source = CachedBase(base_bytes);
let (actual, bytes) = decode_frame_with(
&frame,
version,
&mut source,
super::geometry::entry_limits(
source_limits
.map_or(available, |limits| available.min(limits.max_decoded_bytes)),
),
)
.map_err(|error| {
if matches!(error, PackError::PackfileTooLarge) {
ResolveFailure::Other(budget_exceeded())
} else {
ResolveFailure::Corrupt(ServerError::invalid_argument("object hash mismatch"))
}
})?;
if actual != id {
return Err(ResolveFailure::Corrupt(ServerError::invalid_argument(
"object hash mismatch",
)));
}
Ok((Arc::from(bytes), depth))
}
.await;
visiting.remove(&location);
let value = result?;
memo.insert(location, value.clone(), budget)?;
Ok(value)
})
}