From 22b4807b437fbe549c3e1299b0431d27b3cdbed8 Mon Sep 17 00:00:00 2001 From: Xiaoyang Han Date: Wed, 7 Oct 2026 23:40:26 +0800 Subject: [PATCH] perf(mst2): admit OBJECT batches before bounded body reads Validate every fixed path and current size/digest fact before opening any body. Reject per-object and whole-batch caps before raw I/O; share only exact OID/digest/size bodies and fully verify different sources. Read actual raw streams with bounded collection, exact EOF/length/full SHA, cancellation and current lease protection. Preserve frame/END ordering and old source semantics. Remove the replaced full-body OBJECT resolver and add thirteen real HTTP regressions. Native and performance results remain separate. --- src/api/router/snapshot_content.rs | 171 +++-- src/api/router/snapshot_content_tests.rs | 29 + .../router/snapshot_objects_bounded_tests.rs | 645 ++++++++++++++++++ src/ceres/api_service/mod.rs | 13 + src/jupiter/service/git_service.rs | 12 + 5 files changed, 806 insertions(+), 64 deletions(-) create mode 100644 src/api/router/snapshot_objects_bounded_tests.rs diff --git a/src/api/router/snapshot_content.rs b/src/api/router/snapshot_content.rs index 6e419a61..a6d9efc4 100644 --- a/src/api/router/snapshot_content.rs +++ b/src/api/router/snapshot_content.rs @@ -12,6 +12,7 @@ use axum::{ use futures::stream::StreamExt; use serde::Deserialize; use serde_json::json; +use sha2::{Digest, Sha256}; use super::{ abs_view_path, guarded_treeframe_response, internal, mst2_error_response, request::Mst2Bytes, @@ -19,21 +20,11 @@ use super::{ use crate::ceres::snapshot::{ chunks::{ChunkProjection, get_or_project}, error::{SnapshotError, SnapshotErrorCode}, - pages::{ - MetadataWalkOutcome, WalkOutcome, base64_of, fetch_raw_blob, hex_of, resolve_abs, - resolve_abs_metadata, - }, + pages::{MetadataWalkOutcome, base64_of, fetch_raw_blob, hex_of, resolve_abs_metadata}, resolver::FsKind, view::validate_scope_relative_path, }; -/// One file resolved at a fixed path with verified content. -struct ResolvedFile { - digest: [u8; 32], - size: u64, - raw: Vec, -} - pub(super) struct ResolvedFileMetadata { fs_kind: FsKind, oid: String, @@ -173,52 +164,72 @@ pub(super) fn fs_kind_str(k: FsKind) -> &'static str { } } -/// Resolve a scope-relative path against the fixed tree and verify the -/// optional `expected_digest`. Absence/directory/intermediate outcomes stay -/// typed errors, never an empty body. #[allow(clippy::result_large_err)] -async fn resolve_file( +async fn read_object( handler: &T, - root_tree: &git_internal::internal::object::tree::Tree, - scope: &str, + file: &ResolvedFileMetadata, path: &str, - expected_digest: Option<&str>, -) -> Result { - let abs_path = abs_view_path(scope, path); - match resolve_abs(handler, root_tree, &abs_path) +) -> Result, Response> { + if file.size > OBJECT_ITEM_MAX { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::ScopeInvalid, + "object exceeds the 256KiB item cap", + ))); + } + let expected_size = file.size as usize; + let mut input = handler + .get_raw_blob_stream_by_hash(&file.oid) .await - .map_err(mst2_error_response)? - { - WalkOutcome::FoundFile { - raw, size, digest, .. - } => { - if let Some(expected) = expected_digest - && expected != format!("sha256:{}", hex_of(&digest)) - { - return Err(mst2_error_response(SnapshotError::new( - SnapshotErrorCode::DigestMismatch, - format!("{path}: content does not match expected_digest"), - ))); - } - Ok(ResolvedFile { digest, size, raw }) + .map_err(|error| { + let code = match error { + crate::common::errors::MegaError::ObjStorageNotFound(_) => { + SnapshotErrorCode::ObjectUnavailable + } + crate::common::errors::MegaError::ObjStorageInconsistent(_) => { + SnapshotErrorCode::IntegrityError + } + _ => SnapshotErrorCode::Internal, + }; + tracing::warn!(error = %error, "fixed-view blob fetch failed"); + mst2_error_response(SnapshotError::new( + code, + "fixed-view content could not be read", + )) + })?; + let mut raw = Vec::with_capacity(expected_size); + let mut hash = Sha256::new(); + while let Some(part) = input.next().await { + let bytes = part.map_err(|error| { + tracing::warn!(error = %error, "fixed-view blob stream failed"); + mst2_error_response(SnapshotError::new( + SnapshotErrorCode::Internal, + "fixed-view content could not be read", + )) + })?; + // Reject oversized producer chunks before copying or hashing them. + if bytes.len() > expected_size - raw.len() { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::IntegrityError, + "fixed blob length disagrees with its verified size fact", + ))); } - WalkOutcome::FoundDir => Err(mst2_error_response(SnapshotError::new( - SnapshotErrorCode::NotDirectory, - format!("{path} is a directory"), - ))), - WalkOutcome::Absent => Err(mst2_error_response(SnapshotError::new( - SnapshotErrorCode::PathNotFound, - format!("{path} absent in the fixed view"), - ))), - WalkOutcome::NotDirectory { symlink } => Err(mst2_error_response(SnapshotError::new( - if symlink { - SnapshotErrorCode::SymlinkTraversal - } else { - SnapshotErrorCode::NotDirectory - }, - format!("{path}: intermediate component is not a directory"), - ))), + hash.update(&bytes); + raw.extend_from_slice(&bytes); } + if raw.len() != expected_size { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::IntegrityError, + "fixed blob length disagrees with its verified size fact", + ))); + } + let digest: [u8; 32] = hash.finalize().into(); + if digest != file.digest { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::DigestMismatch, + format!("{path}: content does not match expected_digest"), + ))); + } + Ok(raw) } #[derive(Deserialize, Debug)] @@ -312,9 +323,7 @@ pub(super) async fn objects( .map_err(mst2_error_response)? .unwrap_or(crate::ceres::snapshot::frame_stream::Encoding::Identity); - // Verify every member at its fixed path before any 200 is produced. - // Members resolve concurrently: a batch is up to 128 files, and a - // sequential S3 read per member dominated cold-mount time. + // Admit the whole fixed-path batch before opening any object body. let handler = state .api_handler(std::path::Path::new("/")) .await @@ -329,7 +338,7 @@ pub(super) async fn objects( let resolved: Vec> = futures::stream::iter(req.items.clone()) .map(move |item| async move { validate_scope_relative_path(&item.path).map_err(mst2_error_response)?; - let f = resolve_file( + let f = resolve_file_metadata( handler_ref, root_ref, scope_ref, @@ -351,21 +360,55 @@ pub(super) async fn objects( .buffered(16) .collect() .await; - let mut unique: Vec<([u8; 32], Vec)> = Vec::new(); - let mut seen: Vec<[u8; 32]> = Vec::new(); - let mut logical_bytes = 0u64; + let mut sources: Vec<(ObjectItem, ResolvedFileMetadata)> = Vec::new(); + let mut content_sizes: Vec<([u8; 32], u64)> = Vec::new(); + let mut planned_bytes = 0u64; for pair in resolved { - let (_, f) = pair?; - if !seen.contains(&f.digest) { - if logical_bytes as usize + f.raw.len() > OBJECT_TOTAL_MAX { + let (item, f) = pair?; + if let Some((_, size)) = content_sizes.iter().find(|(digest, _)| *digest == f.digest) { + if *size != f.size { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::IntegrityError, + "fixed content digest has conflicting verified sizes", + ))); + } + } else { + if planned_bytes + f.size > OBJECT_TOTAL_MAX as u64 { return Err(mst2_error_response(SnapshotError::new( SnapshotErrorCode::ScopeInvalid, "unique object content exceeds the 8MiB batch cap", ))); } - logical_bytes += f.raw.len() as u64; - seen.push(f.digest); - unique.push((f.digest, f.raw)); + planned_bytes += f.size; + content_sizes.push((f.digest, f.size)); + } + if let Some((_, source)) = sources.iter().find(|(_, source)| source.oid == f.oid) { + if source.digest != f.digest || source.size != f.size { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::IntegrityError, + "fixed object has conflicting verified facts", + ))); + } + } else { + sources.push((item, f)); + } + } + let loaded = futures::stream::iter(sources) + .map(|(item, file)| async move { + let raw = read_object(handler_ref, &file, &item.path).await?; + Ok::<_, Response>((file.digest, raw)) + }) + .buffered(16); + tokio::pin!(loaded); + let mut unique: Vec<([u8; 32], Vec)> = Vec::new(); + let mut seen: Vec<[u8; 32]> = Vec::new(); + let mut logical_bytes = 0u64; + while let Some(pair) = loaded.next().await { + let (digest, raw) = pair?; + if !seen.contains(&digest) { + logical_bytes += raw.len() as u64; + seen.push(digest); + unique.push((digest, raw)); } } diff --git a/src/api/router/snapshot_content_tests.rs b/src/api/router/snapshot_content_tests.rs index 5399ad42..1e0dcd55 100644 --- a/src/api/router/snapshot_content_tests.rs +++ b/src/api/router/snapshot_content_tests.rs @@ -79,6 +79,9 @@ use crate::{ const TOKEN: &str = "mst2-fixed-content-test"; +#[path = "snapshot_objects_bounded_tests.rs"] +mod bounded_objects; + #[path = "snapshot_session_tests.rs"] mod durable_sessions; @@ -99,6 +102,7 @@ struct ReadCounts { whole: AtomicUsize, range: AtomicUsize, bytes: AtomicUsize, + object_fault: std::sync::Mutex>, } impl ReadCounts { @@ -134,6 +138,11 @@ impl MegaObjectStorage for CountingStorage { async fn get_stream(&self, key: &ObjectKey) -> OrbitResult<(ObjectByteStream, ObjectMeta)> { self.counts.whole.fetch_add(1, Ordering::SeqCst); let (stream, meta) = self.inner.inner.get_stream(key).await?; + let fault = self.counts.object_fault.lock().unwrap().clone(); + let stream = match fault { + Some(fault) if fault.oid == key.key => fault.stream(), + _ => stream, + }; let counts = self.counts.clone(); let stream = stream.map(move |part| { if let Ok(bytes) = &part { @@ -328,6 +337,14 @@ impl Fixture { } async fn new_with_pg_config_and_directories(rebuildable: bool, directory_count: usize) -> Self { + Self::new_with_pg_config_directories_and_objects(rebuildable, directory_count, &[]).await + } + + async fn new_with_pg_config_directories_and_objects( + rebuildable: bool, + directory_count: usize, + objects: &[(String, Vec)], + ) -> Self { let temp = tempfile::tempdir().unwrap(); let mut config = isolated_config(temp.path().join("config")); config.monorepo.push_policy = PushPolicy::Trunk; @@ -414,6 +431,18 @@ impl Fixture { )); extra_trees.push(child); } + for (name, raw) in objects { + let oid = storage + .git_service + .save_object_from_raw(Bytes::copy_from_slice(raw)) + .await + .unwrap(); + project_items.push(item( + TreeItemMode::Blob, + ObjectHash::from_hex_for_kind(HashKind::Sha1, &oid).unwrap(), + name, + )); + } let project = tree(project_items); let old_tip = Commit::from_tree_id_with_kind( HashKind::Sha1, diff --git a/src/api/router/snapshot_objects_bounded_tests.rs b/src/api/router/snapshot_objects_bounded_tests.rs new file mode 100644 index 00000000..6475c53a --- /dev/null +++ b/src/api/router/snapshot_objects_bounded_tests.rs @@ -0,0 +1,645 @@ +use std::io; + +use tokio::{sync::Notify, time::timeout}; + +use super::*; + +#[derive(Clone)] +pub(super) struct StreamFault { + pub(super) oid: String, + kind: FaultKind, +} + +#[derive(Clone)] +enum FaultKind { + Parts(Vec), + LateError(Bytes), + Oversized(Bytes, Arc), + Held { + raw: Bytes, + entered: Arc, + release: Arc, + drops: Arc, + }, +} + +struct DropCount(Arc); + +impl Drop for DropCount { + fn drop(&mut self) { + self.0.fetch_add(1, Ordering::SeqCst); + } +} + +impl StreamFault { + pub(super) fn stream(self) -> ObjectByteStream { + match self.kind { + FaultKind::Parts(parts) => Box::pin(futures::stream::iter(parts.into_iter().map(Ok))), + FaultKind::LateError(raw) => Box::pin(futures::stream::iter([ + Ok(raw), + Err(io::Error::other("test late object stream failure")), + ])), + FaultKind::Oversized(raw, tail_polls) => Box::pin(futures::stream::unfold( + (Some(raw), tail_polls), + |(raw, tail_polls)| async move { + match raw { + Some(raw) => Some((Ok(raw), (None, tail_polls))), + None => { + tail_polls.fetch_add(1, Ordering::SeqCst); + Some(( + Err(io::Error::other("oversized tail must not be polled")), + (None, tail_polls), + )) + } + } + }, + )), + FaultKind::Held { + raw, + entered, + release, + drops, + } => Box::pin(futures::stream::unfold( + (raw, entered, release, DropCount(drops), 0u8), + |(raw, entered, release, drop_count, turn)| async move { + match turn { + 0 => Some((Ok(raw.slice(..1)), (raw, entered, release, drop_count, 1))), + 1 => { + entered.notify_one(); + release.notified().await; + Some((Ok(raw.slice(1..)), (raw, entered, release, drop_count, 2))) + } + _ => None, + } + }, + )), + } + } +} + +fn files(count: usize, size: usize) -> Vec<(String, Vec)> { + let seed = uuid::Uuid::new_v4(); + (0..count) + .map(|index| { + let mut raw = vec![index as u8; size]; + let prefix = format!("blob 3\0abc\0{seed}-{index}"); + raw[..prefix.len()].copy_from_slice(prefix.as_bytes()); + (format!("object-{index:03}"), raw) + }) + .collect() +} + +async fn fixture(objects: &[(String, Vec)]) -> Fixture { + Fixture::new_with_pg_config_directories_and_objects(false, 0, objects).await +} + +fn body(objects: &[(String, Vec)]) -> Body { + Body::from(request_bytes(objects)) +} + +fn request_bytes(objects: &[(String, Vec)]) -> Vec { + serde_json::to_vec(&json!({ + "items":objects.iter().map(|(name, raw)| json!({ + "path":format!("/{name}"), + "expected_digest":format!("sha256:{}", hex_of(&digest(raw))), + })).collect::>(), + "encoding":"identity", + })) + .unwrap() +} + +async fn object_oid(fixture: &Fixture, path: &str) -> String { + let handler = MonoApiService::from(&fixture.state); + let context = fixture + .state + .storage + .mono_storage() + .get_main_ref("/project") + .await + .unwrap() + .unwrap(); + let tree = handler + .get_tree_by_hash(&context.ref_tree_hash) + .await + .unwrap(); + match resolve_abs_metadata(&handler, &tree, path).await.unwrap() { + MetadataWalkOutcome::FoundFile { oid, .. } => oid, + other => panic!("object test path did not resolve: {other:?}"), + } +} + +async fn install_fault(fixture: &Fixture, path: &str, kind: FaultKind) { + let oid = object_oid(fixture, path).await; + *fixture.counts.object_fault.lock().unwrap() = Some(StreamFault { oid, kind }); +} + +async fn assert_objects(response: Response, request: &[u8], objects: &[(String, Vec)]) { + assert_eq!(response.status(), 200); + assert_eq!( + response.headers()["x-mega-request-digest"], + format!("sha256:{}", hex_of(&digest(request))) + ); + let encoded = to_bytes(response.into_body(), 10 * 1024 * 1024) + .await + .unwrap(); + let frames = parse_stream(&encoded).unwrap(); + let mut actual = Vec::new(); + let mut end = None; + for frame in frames { + match frame { + Frame::Object(payload) => { + assert!(end.is_none(), "END must be terminal"); + actual.extend(payload.objects); + } + Frame::End(record) => { + assert!(end.is_none()); + end = Some(record); + } + other => panic!("unexpected OBJECT frame: {other:?}"), + } + } + let mut expected = Vec::new(); + for (_, raw) in objects { + let id = digest(raw); + if !expected.iter().any(|(digest, _)| *digest == id) { + expected.push((id, raw.clone())); + } + } + assert_eq!(actual, expected); + let end = end.unwrap(); + assert_eq!(end.request_item_count, objects.len() as u32); + assert_eq!(end.unique_unit_count, expected.len() as u32); + assert_eq!( + end.logical_bytes, + expected + .iter() + .map(|(_, bytes)| bytes.len() as u64) + .sum::() + ); + assert_eq!(end.request_body_sha256, digest(request)); +} + +#[tokio::test] +async fn oversized_item_and_later_invalid_path_reject_entire_batch_before_body_io() { + let fixture = Fixture::new().await; + let link_digest = format!("sha256:{}", hex_of(&digest(b"file"))); + for last in [ + json!({"path":"/file","expected_digest":fixture.digest_string()}), + json!({"path":"/missing","expected_digest":link_digest}), + ] { + let oversized = last["path"] == "/file"; + let request = json!({"items":[ + {"path":"/link","expected_digest":link_digest},last, + ]}); + error( + fixture + .send("POST", "objects", Body::from(request.to_string())) + .await, + if oversized { 400 } else { 404 }, + if oversized { + "SCOPE_INVALID" + } else { + "PATH_NOT_FOUND" + }, + false, + ) + .await; + fixture.counts.assert(0, 0); + } +} + +#[tokio::test] +async fn unique_batch_over_eight_mib_rejects_before_any_body_io() { + let objects = files(33, 256 * 1024); + let fixture = fixture(&objects).await; + error( + fixture.send("POST", "objects", body(&objects)).await, + 400, + "SCOPE_INVALID", + false, + ) + .await; + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn cap_boundaries_and_empty_object_keep_exact_raw_bytes_and_end_counts() { + let mut objects = files(32, 256 * 1024); + objects.push(("empty".into(), Vec::new())); + let fixture = fixture(&objects[..32]).await; + let request = request_bytes(&objects); + assert_objects( + fixture + .send("POST", "objects", Body::from(request.clone())) + .await, + &request, + &objects, + ) + .await; + fixture.counts.assert(33, 8 * 1024 * 1024); +} + +#[tokio::test] +async fn every_alias_is_admitted_and_exact_oid_body_is_loaded_once() { + let raw = files(1, 8192).remove(0).1; + let objects: Vec<_> = (0..128) + .map(|i| (format!("alias-{i:03}"), raw.clone())) + .collect(); + let fixture = fixture(&objects).await; + let request = request_bytes(&objects); + assert_objects( + fixture + .send("POST", "objects", Body::from(request.clone())) + .await, + &request, + &objects, + ) + .await; + fixture.counts.assert(1, raw.len()); + fixture.counts.reset(); + let mut request: Value = serde_json::from_slice(&request).unwrap(); + request["items"][127]["expected_digest"] = json!(format!("sha256:{}", "00".repeat(32))); + error( + fixture + .send("POST", "objects", Body::from(request.to_string())) + .await, + 409, + "DIGEST_MISMATCH", + false, + ) + .await; + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn conflicting_sizes_reject_before_io_and_distinct_oids_still_verify_each_body() { + let objects = files(2, 8192); + let fixture = fixture(&objects).await; + let oid = object_oid(&fixture, "/object-001").await; + let db = fixture + .state + .storage + .mono_storage() + .get_connection() + .clone(); + let fact = mst2_verified_object::Entity::find() + .filter(mst2_verified_object::Column::GitOid.eq(oid)) + .one(&db) + .await + .unwrap() + .unwrap(); + let mut fact = fact.into_active_model(); + fact.raw_sha256 = Set(digest(&objects[0].1).to_vec()); + fact.size = Set(8193); + let fact = fact.update(&db).await.unwrap(); + let request = json!({"items":[ + {"path":"/object-000","expected_digest":format!("sha256:{}", hex_of(&digest(&objects[0].1)))}, + {"path":"/object-001","expected_digest":format!("sha256:{}", hex_of(&digest(&objects[0].1)))}, + ]}); + error( + fixture + .send("POST", "objects", Body::from(request.to_string())) + .await, + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + fixture.counts.assert(0, 0); + let mut fact = fact.into_active_model(); + fact.size = Set(8192); + fact.update(&db).await.unwrap(); + error( + fixture + .send("POST", "objects", Body::from(request.to_string())) + .await, + 409, + "DIGEST_MISMATCH", + false, + ) + .await; + fixture.counts.assert(2, 16384); +} + +#[tokio::test] +async fn missing_or_invalid_current_facts_fail_before_body_io() { + let fixture = Fixture::new().await; + let request = json!({"items":[{"path":"/file","expected_digest":fixture.digest_string()}]}); + let fact = fixture.fact().await; + fixture.delete_fact().await; + error( + fixture + .send("POST", "objects", Body::from(request.to_string())) + .await, + 503, + "METADATA_NOT_READY", + true, + ) + .await; + fixture.counts.assert(0, 0); + let mut invalid = fact; + invalid.raw_sha256 = vec![0; 31]; + fixture.replace_fact(invalid).await; + error( + fixture + .send("POST", "objects", Body::from(request.to_string())) + .await, + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn oversized_stream_chunk_is_rejected_without_copying_or_polling_tail() { + let objects = files(1, 8192); + let fixture = fixture(&objects).await; + let tail_polls = Arc::new(AtomicUsize::new(0)); + install_fault( + &fixture, + "/object-000", + FaultKind::Oversized(Bytes::from(vec![7; 8193]), tail_polls.clone()), + ) + .await; + error( + fixture.send("POST", "objects", body(&objects)).await, + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + fixture.counts.assert(1, 8193); + assert_eq!(tail_polls.load(Ordering::SeqCst), 0); +} + +#[tokio::test] +async fn truncated_wrong_sha_and_late_stream_error_never_produce_200() { + let objects = files(1, 8192); + let fixture = fixture(&objects).await; + for (fault, status, code, bytes) in [ + ( + FaultKind::Parts(vec![Bytes::copy_from_slice(&objects[0].1[..8191])]), + 502, + "INTEGRITY_ERROR", + 8191, + ), + ( + FaultKind::Parts(vec![Bytes::from(vec![0; 8192])]), + 409, + "DIGEST_MISMATCH", + 8192, + ), + ( + FaultKind::LateError(Bytes::copy_from_slice(&objects[0].1)), + 500, + "INTERNAL", + 8192, + ), + ( + FaultKind::Parts(vec![ + Bytes::copy_from_slice(&objects[0].1), + Bytes::from_static(b"x"), + ]), + 502, + "INTEGRITY_ERROR", + 8193, + ), + ] { + fixture.counts.reset(); + install_fault(&fixture, "/object-000", fault).await; + error( + fixture.send("POST", "objects", body(&objects)).await, + status, + code, + code == "INTERNAL", + ) + .await; + fixture.counts.assert(1, bytes); + } +} + +#[tokio::test] +async fn a_later_object_stream_failure_keeps_earlier_verified_data_unpublished() { + let objects = files(2, 8192); + let fixture = fixture(&objects).await; + install_fault( + &fixture, + "/object-001", + FaultKind::LateError(Bytes::copy_from_slice(&objects[1].1)), + ) + .await; + error( + fixture.send("POST", "objects", body(&objects)).await, + 500, + "INTERNAL", + true, + ) + .await; + fixture.counts.assert(2, 16384); +} + +#[tokio::test] +async fn missing_actual_object_body_keeps_source_error_classification_and_can_retry() { + let objects = files(1, 8192); + let fixture = fixture(&objects).await; + let oid = object_oid(&fixture, "/object-000").await; + fixture + .state + .storage + .git_service + .obj_storage + .inner + .delete(&ObjectKey { + namespace: ObjectNamespace::Git, + key: oid.clone(), + }) + .await + .unwrap(); + error( + fixture.send("POST", "objects", body(&objects)).await, + 503, + "OBJECT_UNAVAILABLE", + false, + ) + .await; + fixture.counts.assert(1, 0); + fixture + .state + .storage + .git_service + .save_object_from_model(objects[0].1.clone(), &oid) + .await + .unwrap(); + fixture.counts.reset(); + let request = request_bytes(&objects); + assert_objects( + fixture + .send("POST", "objects", Body::from(request.clone())) + .await, + &request, + &objects, + ) + .await; + fixture.counts.assert(1, 8192); +} + +#[tokio::test] +async fn old_snapshot_objects_keep_fixed_oid_after_real_publication_advances() { + let objects = files(1, 8192); + let fixture = fixture(&objects).await; + let mono = fixture.state.storage.mono_storage(); + let old = mono.get_main_ref("/project").await.unwrap().unwrap(); + let old_commit = ObjectHash::from_hex_for_kind(HashKind::Sha1, &old.ref_commit_hash).unwrap(); + let oid = object_oid(&fixture, "/object-000").await; + let next_tree = tree(vec![item( + TreeItemMode::Blob, + ObjectHash::from_hex_for_kind(HashKind::Sha1, &oid).unwrap(), + "new-only", + )]); + let next_commit = Commit::from_tree_id_with_kind( + HashKind::Sha1, + next_tree.id, + vec![old_commit], + "OBJECT old snapshot must retain its fixed path", + ) + .unwrap(); + mono.save_mega_trees(vec![next_tree], next_commit.id, None) + .await + .unwrap(); + mono.save_mega_commits(vec![next_commit.clone()], None) + .await + .unwrap(); + publish_native_push(&fixture.state.storage, "/project", old_commit, &next_commit).await; + let head = mono + .read_native_publication_head( + fixture + .state + .storage + .config() + .mst2 + .instance_uuid + .as_deref() + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(head.token.sequence, 2); + assert!(head.token.certificate.is_some()); + assert_ne!( + mono.get_main_ref("/project") + .await + .unwrap() + .unwrap() + .ref_tree_hash, + old.ref_tree_hash + ); + fixture.counts.reset(); + let request = request_bytes(&objects); + assert_objects( + fixture + .send("POST", "objects", Body::from(request.clone())) + .await, + &request, + &objects, + ) + .await; + fixture.counts.assert(1, 8192); +} + +#[tokio::test] +async fn cancelling_actual_object_request_drops_held_stream_and_retry_succeeds() { + let objects = files(1, 8192); + let fixture = fixture(&objects).await; + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let drops = Arc::new(AtomicUsize::new(0)); + install_fault( + &fixture, + "/object-000", + FaultKind::Held { + raw: Bytes::copy_from_slice(&objects[0].1), + entered: entered.clone(), + release, + drops: drops.clone(), + }, + ) + .await; + let task = tokio::spawn(fixture.app.clone().oneshot(fixture.request( + "POST", + "objects", + body(&objects), + ))); + timeout(Duration::from_secs(10), entered.notified()) + .await + .unwrap(); + assert!(!task.is_finished()); + fixture.counts.assert(1, 1); + task.abort(); + assert!(task.await.err().unwrap().is_cancelled()); + assert_eq!(drops.load(Ordering::SeqCst), 1); + *fixture.counts.object_fault.lock().unwrap() = None; + fixture.counts.reset(); + let request = request_bytes(&objects); + assert_objects( + fixture + .send("POST", "objects", Body::from(request.clone())) + .await, + &request, + &objects, + ) + .await; + fixture.counts.assert(1, 8192); +} + +#[tokio::test] +async fn lease_revoked_during_object_load_cannot_deliver_verified_frames() { + let objects = files(1, 8192); + let fixture = fixture(&objects).await; + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let drops = Arc::new(AtomicUsize::new(0)); + install_fault( + &fixture, + "/object-000", + FaultKind::Held { + raw: Bytes::copy_from_slice(&objects[0].1), + entered: entered.clone(), + release: release.clone(), + drops: drops.clone(), + }, + ) + .await; + let task = tokio::spawn(fixture.app.clone().oneshot(fixture.request( + "POST", + "objects", + body(&objects), + ))); + timeout(Duration::from_secs(10), entered.notified()) + .await + .unwrap(); + let revoked = fixture + .app + .clone() + .oneshot( + Request::builder() + .method("DELETE") + .uri(format!("/api/v2/snapshots/leases/{}", fixture.lease)) + .header("authorization", format!("Bearer {TOKEN}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(revoked.status(), 200); + release.notify_one(); + let response = timeout(Duration::from_secs(10), task) + .await + .unwrap() + .unwrap() + .unwrap(); + error(response, 410, "LEASE_EXPIRED", false).await; + fixture.counts.assert(1, 8192); + assert_eq!(drops.load(Ordering::SeqCst), 1); +} diff --git a/src/ceres/api_service/mod.rs b/src/ceres/api_service/mod.rs index e5e0f4be..32bbf780 100644 --- a/src/ceres/api_service/mod.rs +++ b/src/ceres/api_service/mod.rs @@ -178,6 +178,19 @@ pub trait ApiHandler: Send + Sync { } } + async fn get_raw_blob_stream_by_hash( + &self, + hash: &str, + ) -> Result { + let storage = self.get_context(); + match storage.git_service.get_object_stream(hash).await { + Ok(stream) => Ok(stream), + Err(error) => Err(storage + .classify_blob_objstorage_not_found(hash, error) + .await), + } + } + /// Preview unified diff for a single file change async fn preview_file_diff( &self, diff --git a/src/jupiter/service/git_service.rs b/src/jupiter/service/git_service.rs index 17cf70f4..bf42b5ee 100644 --- a/src/jupiter/service/git_service.rs +++ b/src/jupiter/service/git_service.rs @@ -107,6 +107,18 @@ impl GitService { Ok(data) } + pub async fn get_object_stream(&self, hash: &str) -> Result { + if !is_full_hex_object_id(hash) { + return Err(MegaError::Other("Invalid object ID format".to_string())); + } + let key = ObjectKey { + namespace: ObjectNamespace::Git, + key: hash.to_string(), + }; + let (stream, _) = self.obj_storage.inner.get_stream(&key).await?; + Ok(stream) + } + pub fn get_objects_stream(&self, hashes: Vec) -> MultiObjectByteStream<'_> { // Filter out obviously invalid object ids early to avoid spurious backend requests. // Callers that need strict validation should validate up-front and return 4xx.