Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
472 changes: 472 additions & 0 deletions src/api/router/snapshot_chunks_bounded_tests.rs

Large diffs are not rendered by default.

209 changes: 175 additions & 34 deletions src/api/router/snapshot_content.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,12 +15,14 @@ use serde_json::json;
use sha2::{Digest, Sha256};

use super::{
abs_view_path, guarded_treeframe_response, internal, mst2_error_response, request::Mst2Bytes,
abs_view_path, guarded_treeframe_response, guarded_treeframe_response_with_budget, internal,
mst2_error_response, request::Mst2Bytes,
};
use crate::ceres::snapshot::{
chunks::{ChunkProjection, get_or_project},
chunks::{ChunkProjection, get_or_project_stream, projection_reservation_bytes},
content_budget::{PROJECTION_LIVE_BYTES, reserve_range_work, reserve_response},
error::{SnapshotError, SnapshotErrorCode},
pages::{MetadataWalkOutcome, base64_of, fetch_raw_blob, hex_of, resolve_abs_metadata},
pages::{MetadataWalkOutcome, base64_of, hex_of, resolve_abs_metadata},
resolver::FsKind,
view::validate_scope_relative_path,
};
Expand Down Expand Up @@ -475,20 +477,24 @@ async fn project_for<T: crate::ceres::api_service::ApiHandler + ?Sized>(
expected_digest: Option<&str>,
) -> Result<std::sync::Arc<ChunkProjection>, Response> {
let f = resolve_file_metadata(handler, root_tree, scope, path, expected_digest).await?;
project_resolved(handler, &f).await
}

#[allow(clippy::result_large_err)]
async fn project_resolved<T: crate::ceres::api_service::ApiHandler + ?Sized>(
handler: &T,
f: &ResolvedFileMetadata,
) -> Result<std::sync::Arc<ChunkProjection>, Response> {
// The first request for a digest builds the projection from the fixed
// Git object; later requests slice the cached representation. A miss
// rebuilds, never errors with "missing chunk".
let digest = f.digest;
let size = f.size;
let projection = get_or_project(digest, || async move {
let raw = fetch_raw_blob(handler, &f.oid).await?;
if raw.len() as u64 != size {
return Err(SnapshotError::new(
SnapshotErrorCode::IntegrityError,
"fixed blob length disagrees with its verified size fact",
));
}
Ok(raw)
let projection = get_or_project_stream(digest, size, || async move {
handler
.get_raw_blob_stream_by_hash(&f.oid)
.await
.map_err(content_read_error)
})
.await
.map_err(|error| {
Expand Down Expand Up @@ -658,9 +664,91 @@ const CHUNKS_TOTAL_MAX: u64 = 128 * 1024 * 1024;

struct Planned {
projection: std::sync::Arc<ChunkProjection>,
oid: String,
index: u64,
}

struct ResolvedChunk {
file: ResolvedFileMetadata,
index: u64,
map_id: String,
}

fn content_read_error(error: crate::common::errors::MegaError) -> SnapshotError {
use crate::common::errors::MegaError;
let code = match &error {
MegaError::ObjStorageNotFound(_) => SnapshotErrorCode::ObjectUnavailable,
MegaError::ObjStorageInconsistent(_) => SnapshotErrorCode::IntegrityError,
MegaError::Io(error) if error.kind() == std::io::ErrorKind::InvalidData => {
SnapshotErrorCode::IntegrityError
}
_ => SnapshotErrorCode::Internal,
};
tracing::warn!(error = %error, "fixed-view content read failed");
SnapshotError::new(code, "fixed-view content could not be read")
}

#[allow(clippy::result_large_err)]
async fn read_chunk_range<T: crate::ceres::api_service::ApiHandler + ?Sized>(
handler: &T,
planned: &Planned,
) -> Result<Vec<u8>, Response> {
let projection = &planned.projection;
let len = projection.map.chunk_len(planned.index).map_err(|error| {
mst2_error_response(internal(format!("invalid admitted chunk: {error}")))
})?;
let start = planned
.index
.checked_mul(mst2_codec::chunkmap::CHUNK_SIZE as u64)
.ok_or_else(|| mst2_error_response(internal("chunk offset overflow")))?;
let end = start
.checked_add(len)
.ok_or_else(|| mst2_error_response(internal("chunk end overflow")))?;
let (mut input, meta) = handler
.get_raw_blob_range_stream_exact(&planned.oid, start, end)
.await
.map_err(|error| mst2_error_response(content_read_error(error)))?
.ok_or_else(|| {
mst2_error_response(SnapshotError::new(
SnapshotErrorCode::RangeNotSupported,
"fixed source does not support exact raw ranges",
))
})?;
if u64::try_from(meta.size).ok() != Some(projection.map.file_size) {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::IntegrityError,
"range source size disagrees with the fixed verified fact",
)));
}
let mut raw = Vec::new();
raw.try_reserve_exact(len as usize).map_err(|_| {
mst2_error_response(SnapshotError::new(
SnapshotErrorCode::TemporaryUnavailable,
"range allocation could not be admitted",
))
})?;
while let Some(part) = input.next().await {
let bytes = part.map_err(|error| {
tracing::warn!(error = %error, "fixed-view range stream failed");
mst2_error_response(SnapshotError::new(
SnapshotErrorCode::ObjectUnavailable,
"fixed-view range stream failed",
))
})?;
if bytes.len() > len as usize - raw.len() {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::IntegrityError,
"range response exceeds its exact requested length",
)));
}
raw.extend_from_slice(&bytes);
}
projection
.verify_chunk(planned.index, &raw)
.map_err(mst2_error_response)?;
Ok(raw)
}

#[allow(clippy::result_large_err)]
pub(super) async fn chunks(
state: State<crate::api::MonoApiServiceState>,
Expand Down Expand Up @@ -693,33 +781,30 @@ pub(super) async fn chunks(
let root_tree = fixed_root_tree(handler.as_ref(), &ctx.root_tree_oid).await?;
let scope = ctx.built.descriptor.scope.clone();
let mut planned: Vec<Planned> = Vec::new();
let mut resolved: Vec<ResolvedChunk> = Vec::new();
let mut units: Vec<(String, u64)> = Vec::new();
let mut distinct: std::collections::HashMap<[u8; 32], u64> = std::collections::HashMap::new();
let mut projection_bytes = 0usize;
let mut logical_bytes = 0u64;
for item in &req.items {
validate_scope_relative_path(&item.path).map_err(mst2_error_response)?;
let index =
parse_decimal_count(&item.chunk_index, "chunk_index").map_err(mst2_error_response)?;
let proj = project_for(
let file = resolve_file_metadata(
handler.as_ref(),
&root_tree,
&scope,
&item.path,
Some(&item.expected_digest),
)
.await?;
let want_map = format!("sha256:{}", hex_of(&proj.map_id));
if want_map != item.map_id {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::DigestMismatch,
format!("{}: map_id does not bind to this file", item.path),
)));
}
if index >= proj.map.chunk_count {
let chunk_count = file.size.div_ceil(mst2_codec::chunkmap::CHUNK_SIZE as u64);
if index >= chunk_count {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::ScopeInvalid,
format!(
"{}: chunk_index {index} >= chunk_count {}",
item.path, proj.map.chunk_count
item.path, chunk_count
),
)));
}
Expand All @@ -731,10 +816,8 @@ pub(super) async fn chunks(
format!("duplicate chunk unit map_id={} index={index}", item.map_id),
)));
}
let len = proj
.map
.chunk_len(index)
.map_err(|e| mst2_error_response(internal(e.to_string())))?;
let start = index * mst2_codec::chunkmap::CHUNK_SIZE as u64;
let len = (file.size - start).min(mst2_codec::chunkmap::CHUNK_SIZE as u64);
if logical_bytes + len > CHUNKS_TOTAL_MAX {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::ScopeInvalid,
Expand All @@ -743,27 +826,77 @@ pub(super) async fn chunks(
}
logical_bytes += len;
units.push(unit);
match distinct.entry(file.digest) {
std::collections::hash_map::Entry::Occupied(entry) => {
if *entry.get() != file.size {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::IntegrityError,
"fixed facts disagree for the same content digest",
)));
}
}
std::collections::hash_map::Entry::Vacant(entry) => {
let bytes = projection_reservation_bytes(file.size).map_err(mst2_error_response)?;
projection_bytes = projection_bytes.checked_add(bytes).ok_or_else(|| {
mst2_error_response(internal("projection batch memory overflow"))
})?;
if projection_bytes > PROJECTION_LIVE_BYTES {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::LimitExceeded,
"chunk batch exceeds its live projection memory budget",
)));
}
entry.insert(file.size);
}
}
resolved.push(ResolvedChunk {
file,
index,
map_id: item.map_id.clone(),
});
}
let response_bytes = usize::try_from(logical_bytes)
.ok()
.and_then(|bytes| bytes.checked_add(req.items.len() * 1024 + 1024))
.ok_or_else(|| mst2_error_response(internal("chunk response memory overflow")))?;
let response_memory = reserve_response(response_bytes).map_err(mst2_error_response)?;
for item in resolved {
let proj = project_resolved(handler.as_ref(), &item.file).await?;
let want_map = format!("sha256:{}", hex_of(&proj.map_id));
if want_map != item.map_id {
return Err(mst2_error_response(SnapshotError::new(
SnapshotErrorCode::DigestMismatch,
"map_id does not bind to the fixed file",
)));
}
planned.push(Planned {
projection: proj,
index,
oid: item.file.oid,
index: item.index,
});
}

use crate::ceres::snapshot::frame_stream::FrameStream;
let mut stream = FrameStream::new(1, encoding);
let mut out: Vec<Vec<u8>> = Vec::new();
for p in planned.iter() {
// Re-verified slice (digest + length) from the staged projection.
let bytes = p
.projection
.chunk_bytes(p.index)
.map_err(mst2_error_response)?;
let _work_memory = reserve_range_work().map_err(mst2_error_response)?;
// The source is the current request's fixed OID, never a cached
// handler/backend/credential from a different scope.
let bytes = if p.projection.has_inline_bytes() {
p.projection
.chunk_bytes(p.index)
.map_err(mst2_error_response)?
.to_vec()
} else {
read_chunk_range(handler.as_ref(), p).await?
};
let frame = stream
.chunk(
p.projection.map_id,
p.projection.map.file_content_id,
p.index,
bytes.to_vec(),
bytes,
)
.map_err(mst2_error_response)?;
out.push(frame);
Expand All @@ -776,7 +909,15 @@ pub(super) async fn chunks(
);
out.push(end);

guarded_treeframe_response(&state, &ctx, &snapshot_id, &body, out).map_err(mst2_error_response)
guarded_treeframe_response_with_budget(
&state,
&ctx,
&snapshot_id,
&body,
out,
Some(response_memory),
)
.map_err(mst2_error_response)
}

/// Strict decimal-string parse for unsigned counts (spec 04 §1: no leading
Expand Down
63 changes: 63 additions & 0 deletions src/api/router/snapshot_content_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,9 @@ const TOKEN: &str = "mst2-fixed-content-test";
#[path = "snapshot_objects_bounded_tests.rs"]
mod bounded_objects;

#[path = "snapshot_chunks_bounded_tests.rs"]
mod bounded_chunks;

#[path = "snapshot_session_tests.rs"]
mod durable_sessions;

Expand All @@ -103,6 +106,7 @@ struct ReadCounts {
range: AtomicUsize,
bytes: AtomicUsize,
object_fault: std::sync::Mutex<Option<bounded_objects::StreamFault>>,
chunk_faults: std::sync::Mutex<Vec<bounded_chunks::ChunkFault>>,
}

impl ReadCounts {
Expand Down Expand Up @@ -137,6 +141,31 @@ impl MegaObjectStorage for CountingStorage {

async fn get_stream(&self, key: &ObjectKey) -> OrbitResult<(ObjectByteStream, ObjectMeta)> {
self.counts.whole.fetch_add(1, Ordering::SeqCst);
let chunk_fault = self
.counts
.chunk_faults
.lock()
.unwrap()
.iter()
.find(|fault| fault.oid == key.key)
.cloned();
if let Some(fault) = chunk_fault {
let counts = self.counts.clone();
let size = fault.size;
let stream = fault.full_stream().map(move |part| {
if let Ok(bytes) = &part {
counts.bytes.fetch_add(bytes.len(), Ordering::SeqCst);
}
part
});
return Ok((
Box::pin(stream),
ObjectMeta {
size: size as i64,
..Default::default()
},
));
}
let (stream, meta) = self.inner.inner.get_stream(key).await?;
let fault = self.counts.object_fault.lock().unwrap().clone();
let stream = match fault {
Expand All @@ -163,6 +192,40 @@ impl MegaObjectStorage for CountingStorage {
self.inner.inner.get_range_stream(key, start, end).await
}

async fn get_range_stream_exact(
&self,
key: &ObjectKey,
start: u64,
end: u64,
) -> OrbitResult<Option<(ObjectByteStream, ObjectMeta)>> {
self.counts.range.fetch_add(1, Ordering::SeqCst);
let fault = self
.counts
.chunk_faults
.lock()
.unwrap()
.iter()
.find(|fault| fault.oid == key.key)
.cloned();
if let Some(fault) = fault {
let range = fault.range_stream(start, end)?;
return Ok(range.map(|(stream, meta)| {
let counts = self.counts.clone();
let stream = stream.map(move |part| {
if let Ok(bytes) = &part {
counts.bytes.fetch_add(bytes.len(), Ordering::SeqCst);
}
part
});
(Box::pin(stream) as ObjectByteStream, meta)
}));
}
self.inner
.inner
.get_range_stream_exact(key, start, end)
.await
}

async fn exists(&self, key: &ObjectKey) -> OrbitResult<bool> {
self.inner.inner.exists(key).await
}
Expand Down
Loading
Loading