From 09aadd1ef0db24c4d05c5154c26f910258ca50b0 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 21 Sep 2026 16:26:26 +0200 Subject: [PATCH 01/18] feat(tiered): Support resumable uploads Route resumable sessions to long-term storage and publish completed revisions through high-volume tombstones. --- objectstore-service/docs/architecture.md | 2 + objectstore-service/src/backend/tiered.rs | 373 +++++++++++++++++++++- 2 files changed, 361 insertions(+), 14 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 38d4abe1..9de5fce0 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -81,6 +81,8 @@ and large, infrequently-accessed ones: The threshold is **1 MiB**. `TieredStorage` routes objects at or below this size to the high-volume backend; objects exceeding it go to the long-term backend. +Resumable uploads currently always go to the long-term backend, regardless +of size. See [`backend::StorageConfig`] for available backend implementations. diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 586a1719..af270a42 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -17,7 +17,7 @@ //! //! ## Revision Keys //! -//! Every large-object write stores its payload at a **revision key** in the +//! Every long-term write stores its payload at a **revision key** in the //! long-term backend: `{original_key}/{uuid}`. The UUID suffix is random (no //! monotonicity is guaranteed), so each write targets a distinct LT path //! regardless of whether another write to the same logical key is in progress. @@ -99,18 +99,11 @@ //! //! # Resumable Uploads //! -//! TODO: Update this section when tiered storage implements resumable uploads. -//! -//! Not implemented here yet, so [`TieredStorage`] inherits the unsupported defaults from -//! [`Backend`] and every session creation returns [`ErrorKind::Unsupported`]. A resumable upload -//! will be a regular -//! long-term write whose payload arrives across several requests, reusing the revision keys, -//! changelog phases and compare-and-write commit described above: session creation decides -//! the tier from the declared total length and returns [`ErrorKind::Unsupported`] if that tier -//! cannot support it, -//! non-final chunks pass straight through to the upstream session, and the final chunk runs -//! the long-term write sequence. +//! Resumable uploads are always written to the long-term backend, regardless of size. +//! A Resumable Upload remains inaccessible until the long-term backend completes and +//! a high-volume tombstone is committed. +use std::num::NonZeroU64; use std::sync::Arc; use std::sync::atomic::Ordering; use std::time::{Duration, SystemTime}; @@ -120,6 +113,7 @@ use bytes::Bytes; use futures_util::StreamExt; use objectstore_types::metadata::Metadata; use objectstore_types::range::ByteRange; +use objectstore_types::resumable::UploadProgress; use objectstore_types::time::Timestamp; use sentry::{Hub, SentryFutureExt}; use serde::{Deserialize, Serialize}; @@ -137,6 +131,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; +use crate::resumable::BackendToken; use crate::stream::{ClientStream, SizedPeek, counting_stream}; /// The threshold up until which we will go to the "high volume" backend. @@ -207,6 +202,7 @@ pub struct TieredStorageConfig { /// and writes of small objects (e.g. BigTable). /// - Objects **> 1 MiB** go to the `long_term` backend — optimized for cost-efficient /// storage of large objects (e.g. GCS). +/// - Resumable uploads always go to the `long_term` backend, regardless of size. /// /// # Redirect Tombstones /// @@ -392,6 +388,37 @@ impl TieredStorage { } } +#[derive(Debug, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +enum ResumableBackendChoice { + LongTerm(String), +} + +#[derive(Debug, Serialize, Deserialize)] +struct TieredResumableToken { + backend: ResumableBackendChoice, + backend_token: BackendToken, + total_length: u64, + time_expires: Option, +} + +impl TieredResumableToken { + fn decode(id: &ObjectId, token: &BackendToken) -> Result { + let session: Self = + serde_json::from_str(token).map_err(|_| ErrorKind::UnknownUploadSession)?; + let ResumableBackendChoice::LongTerm(key) = &session.backend; + let suffix = key + .strip_prefix(id.key()) + .and_then(|key| key.strip_prefix('/')) + .ok_or(ErrorKind::UnknownUploadSession)?; + if NonZeroU64::new(session.total_length).is_none() || uuid::Uuid::parse_str(suffix).is_err() + { + return Err(ErrorKind::UnknownUploadSession.into()); + } + Ok(session) + } +} + #[async_trait::async_trait] impl Backend for TieredStorage { fn name(&self) -> &'static str { @@ -402,6 +429,189 @@ impl Backend for TieredStorage { Ok(self) } + #[tracing::instrument(level = "debug", fields(?id, total_length), skip_all)] + async fn create_upload_session( + &self, + id: &ObjectId, + metadata: &Metadata, + total_length: NonZeroU64, + ) -> Result> { + let revision = new_long_term_revision(id); + let Some(backend_token) = self + .inner + .long_term + .create_upload_session(&revision, metadata, total_length) + .await? + else { + return Ok(None); + }; + let token = TieredResumableToken { + backend: ResumableBackendChoice::LongTerm(revision.key), + backend_token, + total_length: total_length.get(), + time_expires: metadata.time_expires, + }; + Ok(Some(serde_json::to_string(&token).context( + ErrorKind::Internal, + "encoding tiered resumable session", + )?)) + } + + #[tracing::instrument(level = "debug", fields(?id, offset, content_length), skip_all)] + async fn put_chunk( + &self, + id: &ObjectId, + token: &BackendToken, + offset: u64, + content_length: u64, + stream: ClientStream, + ) -> Result { + let session = TieredResumableToken::decode(id, token)?; + let ResumableBackendChoice::LongTerm(key) = &session.backend; + let revision = ObjectId { + context: id.context.clone(), + key: key.clone(), + }; + let end = offset + .checked_add(content_length) + .filter(|end| *end <= session.total_length) + .ok_or(ErrorKind::ChunkExceedsUploadLength { + offset, + content_length, + upload_length: session.total_length, + })?; + + // Non-final request; just forward the chunk. + if end != session.total_length { + let progress = self + .inner + .long_term + .put_chunk( + &revision, + &session.backend_token, + offset, + content_length, + stream, + ) + .await?; + return match progress { + UploadProgress::Incomplete { .. } => Ok(progress), + UploadProgress::Complete => Err(ErrorKind::UploadSessionGone.into()), + }; + } + + // Verify that this is indeed the final request before proceeding + // (i.e. that the client is not sending us a premature offset). + // GCS returns `Complete` when the upload was completed successfully, even if the underlying + // blob has been subsequently deleted, so the intention here is to prevent the creation of a + // dangling tombstone in such cases. + match self.upload_offset(id, token).await? { + UploadProgress::Complete => return Ok(UploadProgress::Complete), + UploadProgress::Incomplete { offset: actual } if actual < offset => { + return Err(ErrorKind::UploadOffsetMismatch { offset: actual }.into()); + } + UploadProgress::Incomplete { .. } => {} + } + + // 1. Read the current HV revision before the final chunk. + let current = match self + .inner + .high_volume + .get_tiered_metadata(id, Timestamp::now()) + .await? + { + TieredMetadata::Tombstone(t) if t.target == revision => { + return Ok(UploadProgress::Complete); + } + TieredMetadata::Tombstone(t) => Some(t.target), + _ => None, + }; + + // FIXME(consistency): It's possible that 2 concurrent `put_chunk` requests A and B + // reach this point simultaneously. + // If a PUT/DELETE on this key is executed between A and B, B will (attempt to) create a + // dangling tombstone. + // + // FIXME(consistency): The next statement potentially creates an orphan in LT. + // (using `ChangeGuard::Assembling` would not solve the problem, but rather introduce more + // subtle race scenarios). + // + // Other consistency issues may exist within the current implementation. + + // 2. Complete the upload in LT. + let progress = self + .inner + .long_term + .put_chunk( + &revision, + &session.backend_token, + offset, + content_length, + stream, + ) + .await?; + if progress != UploadProgress::Complete { + return Ok(progress); + } + + // 3. Track cleanup of the old revision only. + // Another attempt may still publish the new revision, so a failed attempt must not delete it. + let mut guard = self + .record_change(Change { + id: id.clone(), + new: None, + old: current.clone(), + cleanup_after: None, + }) + .await?; + guard.advance(ChangePhase::Written); + + // 4. Publish the tombstone. + let written = self + .inner + .high_volume + .compare_and_write( + id, + current.as_ref(), + TieredWrite::Tombstone(Tombstone { + target: revision, + time_expires: session.time_expires, + }), + Timestamp::now(), + ) + .await?; + guard.advance(ChangePhase::compare_and_write(written)); + Ok(UploadProgress::Complete) + } + + #[tracing::instrument(level = "debug", fields(?id), skip_all)] + async fn upload_offset(&self, id: &ObjectId, token: &BackendToken) -> Result { + let session = TieredResumableToken::decode(id, token)?; + let ResumableBackendChoice::LongTerm(key) = &session.backend; + let revision = ObjectId { + context: id.context.clone(), + key: key.clone(), + }; + self.inner + .long_term + .upload_offset(&revision, &session.backend_token) + .await + } + + #[tracing::instrument(level = "debug", fields(?id), skip_all)] + async fn cancel_upload(&self, id: &ObjectId, token: &BackendToken) -> Result<()> { + let session = TieredResumableToken::decode(id, token)?; + let ResumableBackendChoice::LongTerm(key) = &session.backend; + let revision = ObjectId { + context: id.context.clone(), + key: key.clone(), + }; + self.inner + .long_term + .cancel_upload(&revision, &session.backend_token) + .await + } + #[tracing::instrument(level = "debug", fields(?id), skip_all)] async fn put_object( &self, @@ -998,7 +1208,7 @@ impl MultipartUploadBackend for TieredStorage { #[cfg(test)] mod tests { - use std::num::NonZeroU32; + use std::num::{NonZeroU32, NonZeroU64}; use std::sync::Mutex as StdMutex; use futures::lock::Mutex; @@ -1006,9 +1216,12 @@ mod tests { use objectstore_types::scope::{Scope, Scopes}; use super::*; + use crate::backend::bigtable::{BigTableBackend, BigTableConfig}; use crate::backend::changelog::{InMemoryChangeLog, NoopChangeLog}; + use crate::backend::gcs::{GcsBackend, GcsConfig}; use crate::backend::in_memory::InMemoryBackend; use crate::backend::testing::{Hooks, TestBackend}; + use crate::change_stream::ChangeStreamFactory; use crate::error::Error; use crate::id::ObjectContext; @@ -1042,6 +1255,138 @@ mod tests { (storage, hv, lt, changelog) } + async fn resumable_token( + storage: &TieredStorage, + id: &ObjectId, + metadata: &Metadata, + length: u64, + ) -> BackendToken { + storage + .create_upload_session(id, metadata, NonZeroU64::new(length).unwrap()) + .await + .unwrap() + .unwrap() + } + + #[tokio::test] + async fn resumable_inmemory() -> anyhow::Result<()> { + let (storage, _, _, _) = make_tiered_storage(); + let id = make_id("tiered-resumable-inmemory"); + let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + + // A completed upload creates a logical object. + assert_eq!( + storage + .put_chunk(&id, &token, 0, 3, stream::single("abc")) + .await?, + UploadProgress::Complete + ); + // InMemory consumes the session immediately. + assert_eq!( + storage.upload_offset(&id, &token).await.unwrap_err().kind(), + ErrorKind::UnknownUploadSession + ); + let (_, _, body) = storage + .get_object(&id, Timestamp::now(), None) + .await? + .unwrap(); + assert_eq!(stream::read_to_vec(body).await?, b"abc"); + Ok(()) + } + + #[tokio::test] + async fn resumable_bigtable_and_gcs() -> anyhow::Result<()> { + let streams = ChangeStreamFactory::default(); + let lt = GcsBackend::new( + GcsConfig { + endpoint: Some("http://localhost:8087".into()), + bucket: "test-bucket".into(), + cogs: None, + }, + &streams, + ) + .await?; + let hv = BigTableBackend::new( + BigTableConfig { + endpoint: Some("localhost:8086".into()), + project_id: "testing".into(), + instance_name: "objectstore".into(), + table_name: "objectstore".into(), + connections: None, + rpc_timeout: Duration::from_secs(2), + cogs: None, + }, + &streams, + ) + .await?; + let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(NoopChangeLog)); + let id = make_id(&format!("tiered-resumable-{}", uuid::Uuid::now_v7())); + let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + + // A completed upload creates a logical object. + assert_eq!( + storage + .put_chunk(&id, &token, 0, 3, stream::single("abc")) + .await?, + UploadProgress::Complete + ); + // GCS still reports completion. + assert_eq!( + storage.upload_offset(&id, &token).await?, + UploadProgress::Complete + ); + let (_, _, body) = storage + .get_object(&id, Timestamp::now(), None) + .await? + .unwrap(); + assert_eq!(stream::read_to_vec(body).await?, b"abc"); + Ok(()) + } + + #[tokio::test] + async fn resumable_invalid_chunks() { + let (storage, _, _, _) = make_tiered_storage(); + let id = make_id("resumable-invalid"); + let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + + // Overflow and future offsets are rejected. + assert_eq!( + storage + .put_chunk(&id, &token, u64::MAX, 1, stream::single("x")) + .await + .unwrap_err() + .kind(), + ErrorKind::ChunkExceedsUploadLength { + offset: u64::MAX, + content_length: 1, + upload_length: 3 + } + ); + assert_eq!( + storage + .put_chunk(&id, &token, 2, 1, stream::single("x")) + .await + .unwrap_err() + .kind(), + ErrorKind::UploadOffsetMismatch { offset: 0 } + ); + } + + #[tokio::test] + async fn resumable_cancel() { + let (storage, hv, _, _) = make_tiered_storage(); + let id = make_id("resumable-invalid"); + let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + + // Cancellation removes the open long-term session. + storage.cancel_upload(&id, &token).await.unwrap(); + assert_eq!( + storage.upload_offset(&id, &token).await.unwrap_err().kind(), + ErrorKind::UnknownUploadSession + ); + assert!(!hv.contains(&id)); + } + #[derive(Clone, Debug)] struct ExpiryHook { label: &'static str, @@ -1605,7 +1950,7 @@ mod tests { // --- CAS conflicts --- - #[derive(Debug)] + #[derive(Debug, Copy, Clone)] struct CasConflict; #[async_trait::async_trait] From 6690213c62edfac345f948cbe2148d891911d16c Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Mon, 21 Sep 2026 18:00:33 +0200 Subject: [PATCH 02/18] fix(tiered): Decline resumable uploads at or below 1 MiB Route resumable sessions only to long-term storage, where uploads can be resumed. Update the existing tiered resumable tests to use long-term sizes and document the cutoff. --- objectstore-service/docs/architecture.md | 5 ++- objectstore-service/src/backend/tiered.rs | 55 +++++++++++++++++------ 2 files changed, 44 insertions(+), 16 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 9de5fce0..2db6d764 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -81,8 +81,9 @@ and large, infrequently-accessed ones: The threshold is **1 MiB**. `TieredStorage` routes objects at or below this size to the high-volume backend; objects exceeding it go to the long-term backend. -Resumable uploads currently always go to the long-term backend, regardless -of size. +`TieredStorage` accepts resumable uploads only when their declared size exceeds +1 MiB; those uploads go to the long-term backend. It declines uploads at or below +1 MiB, which can instead use a direct PUT to the high-volume backend. See [`backend::StorageConfig`] for available backend implementations. diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index af270a42..b1d79798 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -99,9 +99,9 @@ //! //! # Resumable Uploads //! -//! Resumable uploads are always written to the long-term backend, regardless of size. -//! A Resumable Upload remains inaccessible until the long-term backend completes and -//! a high-volume tombstone is committed. +//! Resumable uploads are accepted only when their declared size exceeds 1 MiB and are +//! written to the long-term backend. A resumable upload remains inaccessible until +//! the long-term backend completes and a high-volume tombstone is committed. use std::num::NonZeroU64; use std::sync::Arc; @@ -202,7 +202,7 @@ pub struct TieredStorageConfig { /// and writes of small objects (e.g. BigTable). /// - Objects **> 1 MiB** go to the `long_term` backend — optimized for cost-efficient /// storage of large objects (e.g. GCS). -/// - Resumable uploads always go to the `long_term` backend, regardless of size. +/// - Resumable uploads at or below 1 MiB are declined; larger uploads go to `long_term`. /// /// # Redirect Tombstones /// @@ -436,6 +436,10 @@ impl Backend for TieredStorage { metadata: &Metadata, total_length: NonZeroU64, ) -> Result> { + if total_length.get() <= BACKEND_SIZE_THRESHOLD as u64 { + return Ok(None); + } + let revision = new_long_term_revision(id); let Some(backend_token) = self .inner @@ -1272,12 +1276,20 @@ mod tests { async fn resumable_inmemory() -> anyhow::Result<()> { let (storage, _, _, _) = make_tiered_storage(); let id = make_id("tiered-resumable-inmemory"); - let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1]; + let token = + resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await; // A completed upload creates a logical object. assert_eq!( storage - .put_chunk(&id, &token, 0, 3, stream::single("abc")) + .put_chunk( + &id, + &token, + 0, + payload.len() as u64, + stream::single(payload.clone()) + ) .await?, UploadProgress::Complete ); @@ -1290,7 +1302,7 @@ mod tests { .get_object(&id, Timestamp::now(), None) .await? .unwrap(); - assert_eq!(stream::read_to_vec(body).await?, b"abc"); + assert_eq!(stream::read_to_vec(body).await?, payload); Ok(()) } @@ -1321,12 +1333,20 @@ mod tests { .await?; let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(NoopChangeLog)); let id = make_id(&format!("tiered-resumable-{}", uuid::Uuid::now_v7())); - let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1]; + let token = + resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await; // A completed upload creates a logical object. assert_eq!( storage - .put_chunk(&id, &token, 0, 3, stream::single("abc")) + .put_chunk( + &id, + &token, + 0, + payload.len() as u64, + stream::single(payload.clone()) + ) .await?, UploadProgress::Complete ); @@ -1339,7 +1359,7 @@ mod tests { .get_object(&id, Timestamp::now(), None) .await? .unwrap(); - assert_eq!(stream::read_to_vec(body).await?, b"abc"); + assert_eq!(stream::read_to_vec(body).await?, payload); Ok(()) } @@ -1347,7 +1367,8 @@ mod tests { async fn resumable_invalid_chunks() { let (storage, _, _, _) = make_tiered_storage(); let id = make_id("resumable-invalid"); - let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + let length = BACKEND_SIZE_THRESHOLD as u64 + 1; + let token = resumable_token(&storage, &id, &Metadata::default(), length).await; // Overflow and future offsets are rejected. assert_eq!( @@ -1359,12 +1380,12 @@ mod tests { ErrorKind::ChunkExceedsUploadLength { offset: u64::MAX, content_length: 1, - upload_length: 3 + upload_length: length } ); assert_eq!( storage - .put_chunk(&id, &token, 2, 1, stream::single("x")) + .put_chunk(&id, &token, length - 1, 1, stream::single("x")) .await .unwrap_err() .kind(), @@ -1376,7 +1397,13 @@ mod tests { async fn resumable_cancel() { let (storage, hv, _, _) = make_tiered_storage(); let id = make_id("resumable-invalid"); - let token = resumable_token(&storage, &id, &Metadata::default(), 3).await; + let token = resumable_token( + &storage, + &id, + &Metadata::default(), + BACKEND_SIZE_THRESHOLD as u64 + 1, + ) + .await; // Cancellation removes the open long-term session. storage.cancel_upload(&id, &token).await.unwrap(); From 21e8467bb7b265aaa15d8a70bef8a26d73ada197 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 22 Sep 2026 13:42:01 +0200 Subject: [PATCH 03/18] fix(tiered): Simplify resumable token format Store the long-term revision and backend token directly in the tiered token. Decode its authenticated payload without redundant key-shape validation. --- objectstore-service/src/backend/tiered.rs | 66 +++++++---------------- 1 file changed, 18 insertions(+), 48 deletions(-) diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index b1d79798..9ee50f0e 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -388,34 +388,19 @@ impl TieredStorage { } } -#[derive(Debug, Serialize, Deserialize)] -#[serde(rename_all = "snake_case")] -enum ResumableBackendChoice { - LongTerm(String), -} +type LongTermBackendToken = BackendToken; #[derive(Debug, Serialize, Deserialize)] struct TieredResumableToken { - backend: ResumableBackendChoice, - backend_token: BackendToken, + inner: LongTermBackendToken, + revision: String, total_length: u64, time_expires: Option, } impl TieredResumableToken { - fn decode(id: &ObjectId, token: &BackendToken) -> Result { - let session: Self = - serde_json::from_str(token).map_err(|_| ErrorKind::UnknownUploadSession)?; - let ResumableBackendChoice::LongTerm(key) = &session.backend; - let suffix = key - .strip_prefix(id.key()) - .and_then(|key| key.strip_prefix('/')) - .ok_or(ErrorKind::UnknownUploadSession)?; - if NonZeroU64::new(session.total_length).is_none() || uuid::Uuid::parse_str(suffix).is_err() - { - return Err(ErrorKind::UnknownUploadSession.into()); - } - Ok(session) + fn decode(token: &BackendToken) -> Result { + serde_json::from_str(token).map_err(|_| ErrorKind::UnknownUploadSession.into()) } } @@ -441,7 +426,7 @@ impl Backend for TieredStorage { } let revision = new_long_term_revision(id); - let Some(backend_token) = self + let Some(inner) = self .inner .long_term .create_upload_session(&revision, metadata, total_length) @@ -450,8 +435,8 @@ impl Backend for TieredStorage { return Ok(None); }; let token = TieredResumableToken { - backend: ResumableBackendChoice::LongTerm(revision.key), - backend_token, + revision: revision.key, + inner, total_length: total_length.get(), time_expires: metadata.time_expires, }; @@ -470,11 +455,10 @@ impl Backend for TieredStorage { content_length: u64, stream: ClientStream, ) -> Result { - let session = TieredResumableToken::decode(id, token)?; - let ResumableBackendChoice::LongTerm(key) = &session.backend; + let session = TieredResumableToken::decode(token)?; let revision = ObjectId { context: id.context.clone(), - key: key.clone(), + key: session.revision.clone(), }; let end = offset .checked_add(content_length) @@ -490,13 +474,7 @@ impl Backend for TieredStorage { let progress = self .inner .long_term - .put_chunk( - &revision, - &session.backend_token, - offset, - content_length, - stream, - ) + .put_chunk(&revision, &session.inner, offset, content_length, stream) .await?; return match progress { UploadProgress::Incomplete { .. } => Ok(progress), @@ -546,13 +524,7 @@ impl Backend for TieredStorage { let progress = self .inner .long_term - .put_chunk( - &revision, - &session.backend_token, - offset, - content_length, - stream, - ) + .put_chunk(&revision, &session.inner, offset, content_length, stream) .await?; if progress != UploadProgress::Complete { return Ok(progress); @@ -590,29 +562,27 @@ impl Backend for TieredStorage { #[tracing::instrument(level = "debug", fields(?id), skip_all)] async fn upload_offset(&self, id: &ObjectId, token: &BackendToken) -> Result { - let session = TieredResumableToken::decode(id, token)?; - let ResumableBackendChoice::LongTerm(key) = &session.backend; + let session = TieredResumableToken::decode(token)?; let revision = ObjectId { context: id.context.clone(), - key: key.clone(), + key: session.revision.clone(), }; self.inner .long_term - .upload_offset(&revision, &session.backend_token) + .upload_offset(&revision, &session.inner) .await } #[tracing::instrument(level = "debug", fields(?id), skip_all)] async fn cancel_upload(&self, id: &ObjectId, token: &BackendToken) -> Result<()> { - let session = TieredResumableToken::decode(id, token)?; - let ResumableBackendChoice::LongTerm(key) = &session.backend; + let session = TieredResumableToken::decode(token)?; let revision = ObjectId { context: id.context.clone(), - key: key.clone(), + key: session.revision.clone(), }; self.inner .long_term - .cancel_upload(&revision, &session.backend_token) + .cancel_upload(&revision, &session.inner) .await } From 26fa2ff46504f96726683c292abf0b560de8df1f Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 10:49:26 +0200 Subject: [PATCH 04/18] feat(resumable): Expose upload granularity Report upload granularity at session creation and in progress responses. Reject undersized non-final chunks in GCS and tiered storage, and let the Rust client preflight known chunk lengths. --- clients/rust/README.md | 3 + clients/rust/src/error.rs | 11 ++ clients/rust/src/resumable.rs | 101 ++++++++++++++++-- clients/rust/tests/e2e.rs | 3 + objectstore-server/src/auth/service.rs | 5 + objectstore-server/src/endpoints/common.rs | 3 +- objectstore-server/src/endpoints/mod.rs | 11 +- objectstore-server/src/endpoints/resumable.rs | 98 +++++++++++------ objectstore-server/tests/resumable.rs | 6 +- objectstore-service/docs/architecture.md | 11 +- objectstore-service/src/backend/common.rs | 13 ++- objectstore-service/src/backend/counting.rs | 4 + objectstore-service/src/backend/gcs.rs | 49 ++++++++- objectstore-service/src/backend/testing.rs | 9 ++ objectstore-service/src/backend/tiered.rs | 24 +++++ objectstore-service/src/error.rs | 15 +++ objectstore-service/src/service.rs | 5 + objectstore-types/src/resumable.rs | 14 ++- 18 files changed, 326 insertions(+), 59 deletions(-) diff --git a/clients/rust/README.md b/clients/rust/README.md index 0d814b17..af25910f 100644 --- a/clients/rust/README.md +++ b/clients/rust/README.md @@ -167,6 +167,9 @@ If the request fails midway, it will be possible to resume it from the persisted **Important:** resumable uploads do not automatically compress chunk contents. The `compression` setting only records how the object is encoded; the caller must compress the payload accordingly. The object length and all offsets refer to the bytes after compression. +When manually slicing non-final chunks, use multiples of a positive upload granularity +(`upload.granularity()`). +Shorter chunks are rejected, while a larger unaligned chunk may persist only its aligned prefix. ```rust,no_run #![cfg(feature = "resumable-upload-api")] diff --git a/clients/rust/src/error.rs b/clients/rust/src/error.rs index 18d71a4b..6bad3808 100644 --- a/clients/rust/src/error.rs +++ b/clients/rust/src/error.rs @@ -73,6 +73,17 @@ pub enum Error { #[cfg(feature = "resumable-upload-api")] #[error("resumable upload session is not available")] ResumableUploadUnavailable, + /// A non-final resumable chunk is shorter than the upload granularity. + #[cfg(feature = "resumable-upload-api")] + #[error( + "non-final chunk length {chunk_length} is smaller than upload granularity {upload_granularity}" + )] + ChunkTooSmall { + /// The declared length of the chunk. + chunk_length: u64, + /// The upload granularity in bytes. + upload_granularity: u64, + }, } /// A convenience alias that defaults our [`Error`] type. diff --git a/clients/rust/src/resumable.rs b/clients/rust/src/resumable.rs index 3804e7af..c636ee41 100644 --- a/clients/rust/src/resumable.rs +++ b/clients/rust/src/resumable.rs @@ -10,12 +10,13 @@ use std::borrow::Cow; use std::collections::BTreeMap; use std::fmt; +use std::sync::{Arc, OnceLock}; use bytes::Bytes; use objectstore_types::metadata::Metadata; use objectstore_types::resumable::{ - CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, - UploadOffset, + CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_GRANULARITY, HEADER_UPLOAD_LENGTH, + HEADER_UPLOAD_OFFSET, UploadOffset, }; use reqwest::{Body, Method, Response, StatusCode}; use serde::Serialize; @@ -52,6 +53,8 @@ pub struct ResumableUpload { session: Session, key: ObjectKey, token: SessionToken, + total_length: Option, + granularity: Arc>, } impl Session { @@ -87,6 +90,8 @@ impl Session { session: self.clone(), key: key.into(), token, + total_length: None, + granularity: Arc::new(OnceLock::new()), } } } @@ -102,6 +107,14 @@ impl ResumableUpload { &self.token } + /// Returns this upload's granularity, in bytes. + /// + /// This is `None` for a reconstructed handle until a server response reports the value. + /// Zero means the upload has no granularity. + pub fn granularity(&self) -> Option { + self.granularity.get().copied() + } + /// Builds a request for the server's authoritative upload progress. pub fn progress(&self) -> UploadProgressBuilder { UploadProgressBuilder { @@ -150,6 +163,10 @@ impl ResumableUpload { } } + fn validate_chunk_length(&self, offset: u64, chunk_length: u64) -> crate::Result<()> { + validate_chunk_length(self.granularity(), self.total_length, offset, chunk_length) + } + /// Builds a request to cancel this upload session, discarding any uploaded bytes. pub fn cancel(&self) -> CancelUploadBuilder { CancelUploadBuilder { @@ -167,6 +184,37 @@ impl ResumableUpload { } } +fn validate_chunk_length( + granularity: Option, + total_length: Option, + offset: u64, + chunk_length: u64, +) -> crate::Result<()> { + let Some(granularity) = granularity else { + return Ok(()); + }; + if granularity == 0 || chunk_length == 0 || chunk_length >= granularity { + return Ok(()); + } + + let Some(total_length) = total_length else { + // A reconstructed handle does not know whether this is the final chunk. The server + // retains the authoritative total length and performs the same validation. + return Ok(()); + }; + if !offset + .checked_add(chunk_length) + .is_some_and(|end| end < total_length) + { + return Ok(()); + } + + Err(Error::ChunkTooSmall { + chunk_length, + upload_granularity: granularity, + }) +} + /// A builder for [`Session::create_upload`]. #[derive(Debug)] pub struct CreateResumableUploadBuilder { @@ -267,9 +315,14 @@ impl CreateResumableUploadBuilder { } let response: CreateSessionResponse = response.json().await?; - Ok(Some( - self.session.resume_upload(response.key, response.session), - )) + let upload = ResumableUpload { + session: self.session, + key: response.key, + token: response.session, + total_length: Some(self.total_length), + granularity: Arc::new(OnceLock::from(response.granularity)), + }; + Ok(Some(upload)) } } @@ -293,7 +346,7 @@ impl UploadProgressBuilder { .header(HEADER_UPLOAD_OFFSET, "*") .send() .await?; - parse_progress_response(response).await + parse_progress_response(response, &self.upload).await } } @@ -326,6 +379,8 @@ impl PutChunkBuilder { /// /// Returns [`Error::ResumableUploadUnavailable`] when the session expired, was canceled, or /// could not be found. The upload must be restarted with a new session in that case. + /// Returns [`Error::ChunkTooSmall`] before sending the request when this is known to be a + /// non-final chunk shorter than the upload granularity. /// /// ```rust,ignore /// let offset = match upload.put(offset, chunk).send().await { @@ -336,6 +391,8 @@ impl PutChunkBuilder { /// }; /// ``` pub async fn send(self) -> crate::Result { + self.upload + .validate_chunk_length(self.offset, self.length)?; let response = self .upload .request(Method::PUT)? @@ -344,7 +401,7 @@ impl PutChunkBuilder { .body(self.body) .send() .await?; - parse_progress_response(response).await + parse_progress_response(response, &self.upload).await } } @@ -383,7 +440,19 @@ impl CancelUploadBuilder { } } -async fn parse_progress_response(response: Response) -> crate::Result { +async fn parse_progress_response( + response: Response, + upload: &ResumableUpload, +) -> crate::Result { + if let Some(granularity) = response + .headers() + .get(HEADER_UPLOAD_GRANULARITY) + .and_then(|value| value.to_str().ok()) + .and_then(|value| value.parse().ok()) + { + let _ = upload.granularity.set(granularity); + } + match response.status() { StatusCode::NO_CONTENT | StatusCode::CONFLICT => { let offset = parse_offset(&response); @@ -424,3 +493,19 @@ fn parse_offset(response: &Response) -> Option { UploadOffset::Unknown => None, } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_validate_chunk_length() { + assert!(matches!( + validate_chunk_length(Some(4), Some(10), 0, 3), + Err(Error::ChunkTooSmall { + chunk_length: 3, + upload_granularity: 4, + }) + )); + } +} diff --git a/clients/rust/tests/e2e.rs b/clients/rust/tests/e2e.rs index 71b96c88..c59c82f8 100644 --- a/clients/rust/tests/e2e.rs +++ b/clients/rust/tests/e2e.rs @@ -1053,6 +1053,7 @@ async fn test_resumable_upload() { .unwrap() .unwrap(); assert_eq!(upload.key(), "resumable-client"); + assert_eq!(upload.granularity(), Some(0)); assert_eq!( upload.progress().send().await.unwrap(), UploadProgress::Incomplete { offset: 0 } @@ -1064,10 +1065,12 @@ async fn test_resumable_upload() { // Resume the session and finish the upload. let resumed = session.resume_upload(upload.key(), upload.token().clone()); + assert_eq!(resumed.granularity(), None); assert_eq!( resumed.progress().send().await.unwrap(), UploadProgress::Incomplete { offset: 3 } ); + assert_eq!(resumed.granularity(), Some(0)); assert_eq!( resumed.put(0, "bad").send().await.unwrap(), UploadProgress::Incomplete { offset: 3 } diff --git a/objectstore-server/src/auth/service.rs b/objectstore-server/src/auth/service.rs index 32282be0..b16a9ffc 100644 --- a/objectstore-server/src/auth/service.rs +++ b/objectstore-server/src/auth/service.rs @@ -218,6 +218,11 @@ impl AuthAwareService { // --- Resumable upload operations --- + /// Returns the upload granularity currently reported for new sessions, in bytes. + pub fn upload_granularity(&self) -> u64 { + self.service.upload_granularity() + } + /// Auth-aware wrapper around [`StorageService::create_upload_session`]. pub async fn create_upload_session( &self, diff --git a/objectstore-server/src/endpoints/common.rs b/objectstore-server/src/endpoints/common.rs index c59c73bb..45329909 100644 --- a/objectstore-server/src/endpoints/common.rs +++ b/objectstore-server/src/endpoints/common.rs @@ -144,7 +144,8 @@ impl ApiError { ServiceErrorKind::InvalidMetadata | ServiceErrorKind::InvalidUploadId | ServiceErrorKind::ClientStream - | ServiceErrorKind::ChunkExceedsUploadLength { .. } => StatusCode::BAD_REQUEST, + | ServiceErrorKind::ChunkExceedsUploadLength { .. } + | ServiceErrorKind::ChunkTooSmall { .. } => StatusCode::BAD_REQUEST, ServiceErrorKind::UnknownUploadSession => StatusCode::NOT_FOUND, ServiceErrorKind::RangeNotSatisfiable { .. } => StatusCode::RANGE_NOT_SATISFIABLE, ServiceErrorKind::UploadOffsetMismatch { .. } => StatusCode::CONFLICT, diff --git a/objectstore-server/src/endpoints/mod.rs b/objectstore-server/src/endpoints/mod.rs index c86bfa5b..2ade8084 100644 --- a/objectstore-server/src/endpoints/mod.rs +++ b/objectstore-server/src/endpoints/mod.rs @@ -51,8 +51,9 @@ //! //! Session creation requires an `Upload-Length` header carrying the total size of the object //! in bytes, takes the same metadata headers as a regular upload, and requires an empty body. -//! It answers `200 OK` with `{"key", "session"}`; the session field is the token to use in -//! subsequent query parameters. Metadata is fixed at this point and does not change afterwards. +//! It answers `200 OK` with `{"key", "session", "granularity"}`; the session field is the token +//! to use in subsequent query parameters. `granularity` is this upload's persistence unit in bytes, +//! or zero when no unit is imposed. Metadata is fixed at this point and does not change afterwards. //! //! Chunk uploads and offset queries share one request shape, distinguished by the //! `Upload-Offset` header: a byte offset submits the body as the chunk starting there, while @@ -60,9 +61,13 @@ //! `204 No Content` with the authoritative `Upload-Offset` while bytes remain, and //! `201 Created` with `{"key"}` once the upload is complete and the object is available through //! the normal object endpoints. The session is terminal at that point. +//! These responses also carry `Upload-Granularity`, allowing a client reconstructed from a token +//! to learn this upload's granularity. //! The offset in the response may be lower than the end of the last chunk that was sent. //! Backends can e.g. persist only aligned prefixes and discard the remainder, so clients must //! always continue from the returned offset. +//! Non-final chunks shorter than one positive granularity unit are rejected; final chunks are +//! exempt because they persist the remaining bytes. //! Every chunk requires `Content-Length`, even over HTTP/2, while creation and offset queries //! must not carry a request body. //! @@ -73,7 +78,7 @@ //! //! | Status | Meaning | Client action | //! |--------|---------|---------------| -//! | `400` | Malformed session token, missing `Upload-Length`, nonempty offset query, or a chunk exceeding the declared length | Correct the request | +//! | `400` | Malformed session token, missing `Upload-Length`, nonempty offset query, chunk exceeding the declared length, or non-final chunk shorter than the granularity | Correct the request | //! | `404` | The upload session is unknown or does not belong to this object | Start a new session or correct the request | //! | `409` | A chunk's offset does not match the authoritative offset | Query the offset and continue from there | //! | `410` | The session expired or was canceled | Start a new session | diff --git a/objectstore-server/src/endpoints/resumable.rs b/objectstore-server/src/endpoints/resumable.rs index 70815418..1b73b390 100644 --- a/objectstore-server/src/endpoints/resumable.rs +++ b/objectstore-server/src/endpoints/resumable.rs @@ -6,10 +6,10 @@ //! //! | Operation | Request | Success | //! |---|---|---| -//! | Create | `POST /objects/{usecase}/{scopes}/?upload_type=resumable` | `200` + `{"key","session"}` | -//! | Create | `PUT /objects/{usecase}/{scopes}/{key}?upload_type=resumable` | `200` + `{"key","session"}` | -//! | Chunk | `PUT …/{key}?session=` with `Upload-Offset: ` | `204` + `Upload-Offset`, or `201` + `{"key"}` | -//! | Offset query | `PUT …/{key}?session=` with `Upload-Offset: *` | `204` + `Upload-Offset`, or `201` + `{"key"}` | +//! | Create | `POST /objects/{usecase}/{scopes}/?upload_type=resumable` | `200` + `{"key","session","granularity"}` | +//! | Create | `PUT /objects/{usecase}/{scopes}/{key}?upload_type=resumable` | `200` + `{"key","session","granularity"}` | +//! | Chunk | `PUT …/{key}?session=` with `Upload-Offset: ` | `204` + progress headers, or `201` + `{"key"}` | +//! | Offset query | `PUT …/{key}?session=` with `Upload-Offset: *` | `204` + progress headers, or `201` + `{"key"}` | //! | Cancel | `DELETE …/{key}?session=` | `204` | #![expect( @@ -28,8 +28,8 @@ use objectstore_service::id::{ObjectContext, ObjectId}; use objectstore_service::stream::ClientStream; use objectstore_types::metadata::Metadata; use objectstore_types::resumable::{ - CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_OFFSET, UploadOffset, - UploadProgress, + CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_GRANULARITY, HEADER_UPLOAD_OFFSET, + UploadOffset, UploadProgress, }; use objectstore_types::time::Timestamp; @@ -118,6 +118,7 @@ async fn create_session_for_id( let body = Json(CreateSessionResponse { key: id.key().to_owned(), session, + granularity: service.upload_granularity(), }); Ok((StatusCode::OK, body).into_response()) } @@ -134,7 +135,7 @@ async fn create_session_for_id( /// /// Both answer `204 No Content` with the authoritative offset while bytes remain, and /// `201 Created` with the key once the upload is complete, the session is terminal, and the object -/// is available through the normal object endpoints. +/// is available through the normal object endpoints. Both include `Upload-Granularity`. /// The acknowledged offset may be lower than the submitted chunk's end, so clients should continue /// from this response (or from a later explicit offset query), never from local byte accounting alone. pub(super) async fn continue_session( @@ -146,6 +147,7 @@ pub(super) async fn continue_session( MeteredBody(body): MeteredBody, ) -> ApiResult { let key = id.key().to_owned(); + let granularity = service.upload_granularity(); let progress = match offset { UploadOffset::At(offset) => { @@ -170,7 +172,7 @@ pub(super) async fn continue_session( } }; - progress_response(progress, key) + progress_response(progress, key, granularity) } /// Cancels a session, discarding whatever was uploaded. @@ -184,9 +186,20 @@ pub(super) async fn cancel_session( } /// Turns an [`UploadProgress`] outcome into the response shared by chunks and offset queries. -fn progress_response(progress: ApiResult, key: String) -> ApiResult { - let progress = match progress { - Ok(progress) => progress, +fn progress_response( + progress: ApiResult, + key: String, + granularity: u64, +) -> ApiResult { + let mut response = match progress { + Ok(UploadProgress::Incomplete { offset }) => ( + StatusCode::NO_CONTENT, + [(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset))], + ) + .into_response(), + Ok(UploadProgress::Complete) => { + (StatusCode::CREATED, Json(CompleteUploadResponse { key })).into_response() + } Err(ApiError::Service(error)) => match error.kind() { ErrorKind::UploadOffsetMismatch { offset } => { let error = ApiError::Service(error); @@ -194,24 +207,16 @@ fn progress_response(progress: ApiResult, key: String) -> ApiRes response .headers_mut() .insert(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset)); - return Ok(response); + response } _ => return Err(ApiError::Service(error)), }, Err(error) => return Err(error), }; - - let response = match progress { - UploadProgress::Incomplete { offset } => ( - StatusCode::NO_CONTENT, - [(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset))], - ) - .into_response(), - UploadProgress::Complete => { - (StatusCode::CREATED, Json(CompleteUploadResponse { key })).into_response() - } - }; - + response.headers_mut().insert( + HEADER_UPLOAD_GRANULARITY, + http::HeaderValue::from(granularity), + ); Ok(response) } @@ -219,39 +224,51 @@ fn progress_response(progress: ApiResult, key: String) -> ApiRes mod tests { use super::*; - /// Reads a response's status, `Upload-Offset` header, and body. - async fn parts_of(response: Response) -> (StatusCode, Option, String) { + /// Reads a response's status, progress headers, and body. + async fn parts_of(response: Response) -> (StatusCode, Option, Option, String) { let status = response.status(); let offset = response .headers() .get(HEADER_UPLOAD_OFFSET) .map(|v| v.to_str().unwrap().to_owned()); + let granularity = response + .headers() + .get(HEADER_UPLOAD_GRANULARITY) + .map(|v| v.to_str().unwrap().to_owned()); let body = axum::body::to_bytes(response.into_body(), usize::MAX) .await .unwrap(); - (status, offset, String::from_utf8(body.to_vec()).unwrap()) + ( + status, + offset, + granularity, + String::from_utf8(body.to_vec()).unwrap(), + ) } #[tokio::test] async fn incomplete_progress_answers_no_content_with_the_offset() { let progress = Ok(UploadProgress::Incomplete { offset: 262_144 }); - let response = progress_response(progress, "my-key".into()).unwrap(); + let response = progress_response(progress, "my-key".into(), 262_144).unwrap(); - let (status, offset, body) = parts_of(response).await; + let (status, offset, granularity, body) = parts_of(response).await; assert_eq!(status, StatusCode::NO_CONTENT); assert_eq!(offset.as_deref(), Some("262144")); + assert_eq!(granularity.as_deref(), Some("262144")); assert!(body.is_empty(), "204 must not carry a body: {body:?}"); } #[tokio::test] async fn commit_answers_created_with_the_key() { - let response = progress_response(Ok(UploadProgress::Complete), "my-key".into()).unwrap(); + let response = + progress_response(Ok(UploadProgress::Complete), "my-key".into(), 262_144).unwrap(); - let (status, offset, body) = parts_of(response).await; + let (status, offset, granularity, body) = parts_of(response).await; assert_eq!(status, StatusCode::CREATED); assert_eq!(offset, None, "a commit reports no offset"); + assert_eq!(granularity.as_deref(), Some("262144")); assert_eq!(body, r#"{"key":"my-key"}"#); } @@ -259,22 +276,23 @@ mod tests { async fn offset_mismatch_answers_conflict_with_the_authoritative_offset() { let mismatch = ErrorKind::UploadOffsetMismatch { offset: 786_432 }.into(); let response = - progress_response(Err(ApiError::Service(mismatch)), "my-key".into()).unwrap(); + progress_response(Err(ApiError::Service(mismatch)), "my-key".into(), 262_144).unwrap(); - let (status, offset, body) = parts_of(response).await; + let (status, offset, granularity, body) = parts_of(response).await; assert_eq!(status, StatusCode::CONFLICT); assert_eq!( offset.as_deref(), Some("786432"), "the client resynchronizes from this header" ); + assert_eq!(granularity.as_deref(), Some("262144")); assert!(body.contains("786432"), "{body:?}"); } #[tokio::test] async fn other_errors_propagate_unchanged() { let gone = ApiError::Service(ErrorKind::UploadSessionGone.into()); - let error = progress_response(Err(gone), "my-key".into()).unwrap_err(); + let error = progress_response(Err(gone), "my-key".into(), 262_144).unwrap_err(); assert_eq!(error.status(), StatusCode::GONE); let oversized = @@ -283,7 +301,17 @@ mod tests { content_length: 4, upload_length: 10, })); - let error = progress_response(Err(oversized), "my-key".into()).unwrap_err(); + let error = progress_response(Err(oversized), "my-key".into(), 262_144).unwrap_err(); + assert_eq!(error.status(), StatusCode::BAD_REQUEST); + + let too_small = ApiError::Service( + ErrorKind::ChunkTooSmall { + chunk_length: 1, + upload_granularity: 262_144, + } + .into(), + ); + let error = progress_response(Err(too_small), "my-key".into(), 262_144).unwrap_err(); assert_eq!(error.status(), StatusCode::BAD_REQUEST); } } diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index 2cc728ce..461356be 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -8,7 +8,7 @@ use anyhow::Result; use objectstore_server::config::{AuthZ, Config, EncryptionConfig, Service}; use objectstore_test::server::TestServer; use objectstore_types::resumable::{ - CreateSessionResponse, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, + CreateSessionResponse, HEADER_UPLOAD_GRANULARITY, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, }; use reqwest::StatusCode; @@ -87,6 +87,7 @@ async fn create_session_path( .await?; assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; + assert_eq!(created.granularity, 0); Ok(format!( "{OBJECT_PATH}?session={}", created.session.to_base64url() @@ -175,6 +176,7 @@ async fn test_resumable_upload() -> Result<()> { assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; assert_eq!(created.key, "my-key"); + assert_eq!(created.granularity, 0); let session_path = format!("{object}?session={}", created.session.to_base64url()); // The object doesn't exist yet. @@ -190,6 +192,7 @@ async fn test_resumable_upload() -> Result<()> { .await?; assert_eq!(response.status(), StatusCode::NO_CONTENT); assert_eq!(response.headers()[HEADER_UPLOAD_OFFSET], "3"); + assert_eq!(response.headers()[HEADER_UPLOAD_GRANULARITY], "0"); // Rejected: outdated offset. let response = client @@ -218,6 +221,7 @@ async fn test_resumable_upload() -> Result<()> { .await?; assert_eq!(response.status(), StatusCode::NO_CONTENT); assert_eq!(response.headers()[HEADER_UPLOAD_OFFSET], "3"); + assert_eq!(response.headers()[HEADER_UPLOAD_GRANULARITY], "0"); let response = client .put(server.url(&session_path)) .header(HEADER_UPLOAD_OFFSET, "3") diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index c9e64573..1fc52d2e 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -334,10 +334,13 @@ Not all backends support resumable uploads. A backend returns no session when it particular upload; this is a routine outcome rather than an error. Acceptance can depend on the declared size, the metadata, or whether resuming is possible in principle. -The offset returned by every chunk response is authoritative. A backend may accept the full -request body but persist only a prefix (for example, up to an internal alignment boundary), so a -client must not advance by the submitted `Content-Length` on its own. It continues from the -preceding response's offset, or performs an explicit offset query after an ambiguous failure. +The offset returned by every chunk response is authoritative. Each upload reports its granularity +in bytes, or zero when it imposes none. For a positive granularity, a backend may accept +the full request body but persist only the largest aligned prefix, so a client must not advance by +the submitted `Content-Length` on its own. Non-empty, non-final chunks shorter than one granularity +unit are rejected; the final chunk is exempt. Clients should use +granularity-sized multiples and continue from the preceding response's offset, or perform an +explicit offset query after an ambiguous failure. ## Multipart Uploads diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index 8087a914..ea0a2221 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -190,6 +190,16 @@ pub trait Backend: fmt::Debug + Send + Sync + 'static { /// The backend name, used for diagnostics. fn name(&self) -> &'static str; + /// Returns the upload granularity for sessions opened by this backend, in bytes. + /// + /// A value of zero means these uploads have no granularity. A positive value means that a + /// non-final chunk can persist only a multiple of this value. Implementations must reject a + /// non-empty, non-final chunk shorter than one unit with [`ErrorKind::ChunkTooSmall`]. Final + /// chunks and zero-byte offset probes are exempt. + fn upload_granularity(&self) -> u64 { + 0 + } + /// Stores an object at the given path with the given metadata. async fn put_object( &self, @@ -296,7 +306,8 @@ pub trait Backend: fmt::Debug + Send + Sync + 'static { /// /// Returns [`ErrorKind::UnknownUploadSession`] when `token` does not identify an open session, /// and [`ErrorKind::ChunkExceedsUploadLength`] when the chunk would exceed the total length - /// declared when the session was created. + /// declared when the session was created. Returns [`ErrorKind::ChunkTooSmall`] when a non-empty, + /// non-final chunk is shorter than the upload granularity. async fn put_chunk( &self, id: &ObjectId, diff --git a/objectstore-service/src/backend/counting.rs b/objectstore-service/src/backend/counting.rs index 4acf0fcd..36fdb4d3 100644 --- a/objectstore-service/src/backend/counting.rs +++ b/objectstore-service/src/backend/counting.rs @@ -71,6 +71,10 @@ impl Backend for CountingBackend { self.inner.name() } + fn upload_granularity(&self) -> u64 { + self.inner.upload_granularity() + } + async fn put_object( &self, id: &ObjectId, diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index b182eec2..e253f3a3 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -902,6 +902,10 @@ impl Backend for GcsBackend { "gcs" } + fn upload_granularity(&self) -> u64 { + 256 * 1024 + } + fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> { Ok(self) } @@ -1266,6 +1270,14 @@ impl Backend for GcsBackend { content_length, upload_length: session.total_length.get(), })?; + let granularity = self.upload_granularity(); + if content_length > 0 && content_length < granularity && end != session.total_length.get() { + return Err(ErrorKind::ChunkTooSmall { + chunk_length: content_length, + upload_granularity: granularity, + } + .into()); + } let content_range = match content_length { // An empty chunk is equivalent to an offset query. @@ -1965,8 +1977,9 @@ mod tests { async fn test_resumable_empty_chunk_reports_offset_without_writing() -> Result<()> { let backend = create_test_backend().await?; let id = make_id_with_key("resumable-empty-chunk"); + let total_length = 2 * (RESUMABLE_CHUNK_SIZE + 2); let token = backend - .create_upload_session(&id, &Metadata::default(), nonzero(4)) + .create_upload_session(&id, &Metadata::default(), nonzero(total_length as u64)) .await?; // An empty chunk cannot advance a session that still expects bytes. It reports the @@ -1977,19 +1990,34 @@ mod tests { .await?, UploadProgress::Incomplete { offset: 0 } ); + // The backend rejects undersized non-final chunks, so use an unaligned chunk larger + // than one granularity unit to exercise the emulator's alignment behavior. + let chunk_length = RESUMABLE_CHUNK_SIZE + 2; let after_write = backend - .put_chunk(&id, &token, 0, 2, stream::single(b"ab".to_vec())) + .put_chunk( + &id, + &token, + 0, + chunk_length as u64, + stream::single(vec![b'a'; chunk_length]), + ) .await?; assert!(matches!(after_write, UploadProgress::Incomplete { .. })); // Whichever prefix GCS acknowledged, an empty chunk reports that same position rather // than moving it. GCS documents that a chunk "should be a multiple of 256 KiB ... unless // it's the last chunk", and that a client "should not assume that the server received all - // bytes sent in any given request". The emulator acknowledges any length, so the position - // itself is not asserted here. + // bytes sent in any given request". The emulator acknowledges the unaligned tail, so the + // position itself is not asserted here. assert_eq!( backend - .put_chunk(&id, &token, 2, 0, stream::single(Vec::new())) + .put_chunk( + &id, + &token, + chunk_length as u64, + 0, + stream::single(Vec::new()), + ) .await?, after_write ); @@ -2042,6 +2070,17 @@ mod tests { backend.upload_offset(&multi_id, &token).await?, UploadProgress::Incomplete { offset: 0 } ); + let error = backend + .put_chunk(&multi_id, &token, 0, 1, stream::single(b"a".to_vec())) + .await + .unwrap_err(); + assert_eq!( + error.kind(), + ErrorKind::ChunkTooSmall { + chunk_length: 1, + upload_granularity: RESUMABLE_CHUNK_SIZE as u64, + } + ); assert_eq!( backend .put_chunk( diff --git a/objectstore-service/src/backend/testing.rs b/objectstore-service/src/backend/testing.rs index d1eb0feb..7ef00627 100644 --- a/objectstore-service/src/backend/testing.rs +++ b/objectstore-service/src/backend/testing.rs @@ -78,6 +78,11 @@ pub trait Hooks: fmt::Debug + Send + Sync + 'static { "test-backend" } + /// Intercepts [`Backend::upload_granularity`]. Default delegates to `inner`. + fn upload_granularity(&self, inner: &InMemoryBackend) -> u64 { + inner.upload_granularity() + } + /// Intercepts [`Backend::put_object`]. Default delegates to `inner`. async fn put_object( &self, @@ -379,6 +384,10 @@ impl Backend for TestBackend { self.hooks.name() } + fn upload_granularity(&self) -> u64 { + self.hooks.upload_granularity(&self.inner) + } + fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> { Ok(self) } diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 2fec9841..2718ef11 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -410,6 +410,10 @@ impl Backend for TieredStorage { "tiered" } + fn upload_granularity(&self) -> u64 { + self.inner.long_term.upload_granularity() + } + fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> { Ok(self) } @@ -468,6 +472,14 @@ impl Backend for TieredStorage { content_length, upload_length: session.total_length, })?; + let granularity = self.upload_granularity(); + if content_length > 0 && content_length < granularity && end != session.total_length { + return Err(ErrorKind::ChunkTooSmall { + chunk_length: content_length, + upload_granularity: granularity, + } + .into()); + } // Non-final request; just forward the chunk. if end != session.total_length { @@ -1309,6 +1321,18 @@ mod tests { let token = resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await; + let error = storage + .put_chunk(&id, &token, 0, 1, stream::single("a")) + .await + .unwrap_err(); + assert_eq!( + error.kind(), + ErrorKind::ChunkTooSmall { + chunk_length: 1, + upload_granularity: 256 * 1024, + } + ); + // A completed upload creates a logical object. assert_eq!( storage diff --git a/objectstore-service/src/error.rs b/objectstore-service/src/error.rs index b427fc80..b38ed3f6 100644 --- a/objectstore-service/src/error.rs +++ b/objectstore-service/src/error.rs @@ -70,6 +70,13 @@ pub enum ErrorKind { /// The total upload length declared when the session was created. upload_length: u64, }, + /// A non-final resumable chunk is shorter than the upload granularity. + ChunkTooSmall { + /// The declared length of the chunk. + chunk_length: u64, + /// The upload granularity in bytes. + upload_granularity: u64, + }, /// The service cannot accept more work. AtCapacity, /// The requested operation is unsupported. @@ -114,6 +121,13 @@ impl fmt::Display for ErrorKind { f, "chunk at offset {offset} with length {content_length} exceeds upload length {upload_length}" ), + Self::ChunkTooSmall { + chunk_length, + upload_granularity, + } => write!( + f, + "non-final chunk length {chunk_length} is smaller than upload granularity {upload_granularity}" + ), Self::AtCapacity => f.write_str("service at capacity"), Self::Unsupported => f.write_str("unsupported operation"), Self::BackendFailure => f.write_str("backend operation failed"), @@ -192,6 +206,7 @@ impl Error { ErrorKind::UploadSessionGone => Level::DEBUG, ErrorKind::UnknownUploadSession => Level::DEBUG, ErrorKind::ChunkExceedsUploadLength { .. } => Level::DEBUG, + ErrorKind::ChunkTooSmall { .. } => Level::DEBUG, // Indicates that optional functionality is not supported. // We don't want a rogue client spamming us with Sentry errors just by calling an API // that the server doesn't support, so we just log it. diff --git a/objectstore-service/src/service.rs b/objectstore-service/src/service.rs index fb14e608..8a16caee 100644 --- a/objectstore-service/src/service.rs +++ b/objectstore-service/src/service.rs @@ -465,6 +465,11 @@ impl StorageService { // --- Resumable upload operations --- + /// Returns the upload granularity currently reported for new sessions, in bytes. + pub fn upload_granularity(&self) -> u64 { + self.inner.upload_granularity() + } + /// Opens a resumable upload session for an object of `total_length` bytes. /// /// Returns `Ok(None)` for zero-length objects or when the backend declines resumable uploads diff --git a/objectstore-types/src/resumable.rs b/objectstore-types/src/resumable.rs index d1b6ea1d..d71a55ee 100644 --- a/objectstore-types/src/resumable.rs +++ b/objectstore-types/src/resumable.rs @@ -3,6 +3,8 @@ //! A resumable upload writes one object across multiple requests. The client first creates a //! session, declaring the object's complete size with [`HEADER_UPLOAD_LENGTH`]. The server returns //! a [`CreateSessionResponse`] containing an opaque [`SessionToken`] that identifies the upload. +//! The response also reports the upload granularity. Subsequent progress responses +//! repeat that value in [`HEADER_UPLOAD_GRANULARITY`], so reconstructed clients can learn it. //! //! The client then sends chunks with [`HEADER_UPLOAD_OFFSET`] set to the byte position at which //! each chunk starts. If an upload is interrupted, the client can send the wildcard offset @@ -32,6 +34,13 @@ pub const HEADER_UPLOAD_LENGTH: &str = "upload-length"; /// persisted. See [`UploadOffset`]. pub const HEADER_UPLOAD_OFFSET: &str = "upload-offset"; +/// Response header declaring the upload granularity, in bytes. +/// +/// A value of zero means the upload has no granularity. For a positive value, +/// non-final chunks shorter than one granularity unit are rejected, and a backend may persist +/// only a multiple of this value from a larger non-final chunk. +pub const HEADER_UPLOAD_GRANULARITY: &str = "upload-granularity"; + /// The wildcard [`HEADER_UPLOAD_OFFSET`] value that queries the server's offset. const OFFSET_WILDCARD: &str = "*"; @@ -184,6 +193,8 @@ pub struct CreateSessionResponse { pub key: String, /// The opaque session token that identifies the session. pub session: SessionToken, + /// This upload's granularity in bytes. + pub granularity: u64, } /// Response from the request that completes the upload. @@ -205,11 +216,12 @@ mod tests { let response = CreateSessionResponse { key: "key".into(), session: SessionToken::new(b"../opaque +? \xc3\xbc"), + granularity: 262_144, }; assert_eq!( serde_json::to_string(&response)?, - r#"{"key":"key","session":"Li4vb3BhcXVlICs_IMO8"}"# + r#"{"key":"key","session":"Li4vb3BhcXVlICs_IMO8","granularity":262144}"# ); Ok(()) } From 033b50c30625e9232d86e0e4b283b67e15f5bcbf Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 11:37:48 +0200 Subject: [PATCH 05/18] wip --- clients/rust/src/error.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/clients/rust/src/error.rs b/clients/rust/src/error.rs index 6bad3808..a7fd6a49 100644 --- a/clients/rust/src/error.rs +++ b/clients/rust/src/error.rs @@ -73,10 +73,10 @@ pub enum Error { #[cfg(feature = "resumable-upload-api")] #[error("resumable upload session is not available")] ResumableUploadUnavailable, - /// A non-final resumable chunk is shorter than the upload granularity. + /// A non-final chunk is shorter than the upload granularity. #[cfg(feature = "resumable-upload-api")] #[error( - "non-final chunk length {chunk_length} is smaller than upload granularity {upload_granularity}" + "non-final chunk {chunk_length} is smaller than upload granularity {upload_granularity}" )] ChunkTooSmall { /// The declared length of the chunk. From 561239430fed0cad31455e29a35c056e5e8a81cc Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 11:47:42 +0200 Subject: [PATCH 06/18] fix(resumable): Report granularity only at session creation Keep creation-time granularity on new client handles and leave reconstructed handles without it. Remove progress-response granularity metadata and inline the client chunk-length check. --- clients/rust/README.md | 2 +- clients/rust/src/resumable.rs | 103 ++++++------------ clients/rust/tests/e2e.rs | 2 +- objectstore-server/src/endpoints/mod.rs | 2 - objectstore-server/src/endpoints/resumable.rs | 64 ++++------- objectstore-server/tests/resumable.rs | 4 +- objectstore-service/docs/architecture.md | 14 +-- objectstore-types/src/resumable.rs | 10 +- 8 files changed, 62 insertions(+), 139 deletions(-) diff --git a/clients/rust/README.md b/clients/rust/README.md index af25910f..18f40bcc 100644 --- a/clients/rust/README.md +++ b/clients/rust/README.md @@ -168,7 +168,7 @@ If the request fails midway, it will be possible to resume it from the persisted setting only records how the object is encoded; the caller must compress the payload accordingly. The object length and all offsets refer to the bytes after compression. When manually slicing non-final chunks, use multiples of a positive upload granularity -(`upload.granularity()`). +(`upload.granularity()` on a newly created handle). Reconstructed handles do not know this value. Shorter chunks are rejected, while a larger unaligned chunk may persist only its aligned prefix. ```rust,no_run diff --git a/clients/rust/src/resumable.rs b/clients/rust/src/resumable.rs index c636ee41..2e624f41 100644 --- a/clients/rust/src/resumable.rs +++ b/clients/rust/src/resumable.rs @@ -10,13 +10,12 @@ use std::borrow::Cow; use std::collections::BTreeMap; use std::fmt; -use std::sync::{Arc, OnceLock}; use bytes::Bytes; use objectstore_types::metadata::Metadata; use objectstore_types::resumable::{ - CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_GRANULARITY, HEADER_UPLOAD_LENGTH, - HEADER_UPLOAD_OFFSET, UploadOffset, + CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, + UploadOffset, }; use reqwest::{Body, Method, Response, StatusCode}; use serde::Serialize; @@ -54,7 +53,7 @@ pub struct ResumableUpload { key: ObjectKey, token: SessionToken, total_length: Option, - granularity: Arc>, + granularity: Option, } impl Session { @@ -91,7 +90,7 @@ impl Session { key: key.into(), token, total_length: None, - granularity: Arc::new(OnceLock::new()), + granularity: None, } } } @@ -109,10 +108,10 @@ impl ResumableUpload { /// Returns this upload's granularity, in bytes. /// - /// This is `None` for a reconstructed handle until a server response reports the value. + /// This is `None` for a reconstructed handle; only session creation supplies the value. /// Zero means the upload has no granularity. pub fn granularity(&self) -> Option { - self.granularity.get().copied() + self.granularity } /// Builds a request for the server's authoritative upload progress. @@ -164,7 +163,25 @@ impl ResumableUpload { } fn validate_chunk_length(&self, offset: u64, chunk_length: u64) -> crate::Result<()> { - validate_chunk_length(self.granularity(), self.total_length, offset, chunk_length) + let (Some(granularity), Some(total_length)) = (self.granularity, self.total_length) else { + // A reconstructed handle cannot tell whether this is the final chunk. The backend + // retains the authoritative total length and validates the request. + return Ok(()); + }; + if granularity == 0 || chunk_length == 0 || chunk_length >= granularity { + return Ok(()); + } + if !offset + .checked_add(chunk_length) + .is_some_and(|end| end < total_length) + { + return Ok(()); + } + + Err(Error::ChunkTooSmall { + chunk_length, + upload_granularity: granularity, + }) } /// Builds a request to cancel this upload session, discarding any uploaded bytes. @@ -184,37 +201,6 @@ impl ResumableUpload { } } -fn validate_chunk_length( - granularity: Option, - total_length: Option, - offset: u64, - chunk_length: u64, -) -> crate::Result<()> { - let Some(granularity) = granularity else { - return Ok(()); - }; - if granularity == 0 || chunk_length == 0 || chunk_length >= granularity { - return Ok(()); - } - - let Some(total_length) = total_length else { - // A reconstructed handle does not know whether this is the final chunk. The server - // retains the authoritative total length and performs the same validation. - return Ok(()); - }; - if !offset - .checked_add(chunk_length) - .is_some_and(|end| end < total_length) - { - return Ok(()); - } - - Err(Error::ChunkTooSmall { - chunk_length, - upload_granularity: granularity, - }) -} - /// A builder for [`Session::create_upload`]. #[derive(Debug)] pub struct CreateResumableUploadBuilder { @@ -320,7 +306,7 @@ impl CreateResumableUploadBuilder { key: response.key, token: response.session, total_length: Some(self.total_length), - granularity: Arc::new(OnceLock::from(response.granularity)), + granularity: Some(response.granularity), }; Ok(Some(upload)) } @@ -346,7 +332,7 @@ impl UploadProgressBuilder { .header(HEADER_UPLOAD_OFFSET, "*") .send() .await?; - parse_progress_response(response, &self.upload).await + parse_progress_response(response).await } } @@ -380,7 +366,8 @@ impl PutChunkBuilder { /// Returns [`Error::ResumableUploadUnavailable`] when the session expired, was canceled, or /// could not be found. The upload must be restarted with a new session in that case. /// Returns [`Error::ChunkTooSmall`] before sending the request when this is known to be a - /// non-final chunk shorter than the upload granularity. + /// non-final chunk shorter than the upload granularity. Reconstructed handles cannot perform + /// this check because they do not know the upload length or granularity. /// /// ```rust,ignore /// let offset = match upload.put(offset, chunk).send().await { @@ -401,7 +388,7 @@ impl PutChunkBuilder { .body(self.body) .send() .await?; - parse_progress_response(response, &self.upload).await + parse_progress_response(response).await } } @@ -440,19 +427,7 @@ impl CancelUploadBuilder { } } -async fn parse_progress_response( - response: Response, - upload: &ResumableUpload, -) -> crate::Result { - if let Some(granularity) = response - .headers() - .get(HEADER_UPLOAD_GRANULARITY) - .and_then(|value| value.to_str().ok()) - .and_then(|value| value.parse().ok()) - { - let _ = upload.granularity.set(granularity); - } - +async fn parse_progress_response(response: Response) -> crate::Result { match response.status() { StatusCode::NO_CONTENT | StatusCode::CONFLICT => { let offset = parse_offset(&response); @@ -493,19 +468,3 @@ fn parse_offset(response: &Response) -> Option { UploadOffset::Unknown => None, } } - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_validate_chunk_length() { - assert!(matches!( - validate_chunk_length(Some(4), Some(10), 0, 3), - Err(Error::ChunkTooSmall { - chunk_length: 3, - upload_granularity: 4, - }) - )); - } -} diff --git a/clients/rust/tests/e2e.rs b/clients/rust/tests/e2e.rs index c59c82f8..c9234b1b 100644 --- a/clients/rust/tests/e2e.rs +++ b/clients/rust/tests/e2e.rs @@ -1070,7 +1070,7 @@ async fn test_resumable_upload() { resumed.progress().send().await.unwrap(), UploadProgress::Incomplete { offset: 3 } ); - assert_eq!(resumed.granularity(), Some(0)); + assert_eq!(resumed.granularity(), None); assert_eq!( resumed.put(0, "bad").send().await.unwrap(), UploadProgress::Incomplete { offset: 3 } diff --git a/objectstore-server/src/endpoints/mod.rs b/objectstore-server/src/endpoints/mod.rs index 2ade8084..0b1a25d0 100644 --- a/objectstore-server/src/endpoints/mod.rs +++ b/objectstore-server/src/endpoints/mod.rs @@ -61,8 +61,6 @@ //! `204 No Content` with the authoritative `Upload-Offset` while bytes remain, and //! `201 Created` with `{"key"}` once the upload is complete and the object is available through //! the normal object endpoints. The session is terminal at that point. -//! These responses also carry `Upload-Granularity`, allowing a client reconstructed from a token -//! to learn this upload's granularity. //! The offset in the response may be lower than the end of the last chunk that was sent. //! Backends can e.g. persist only aligned prefixes and discard the remainder, so clients must //! always continue from the returned offset. diff --git a/objectstore-server/src/endpoints/resumable.rs b/objectstore-server/src/endpoints/resumable.rs index 1b73b390..8c527233 100644 --- a/objectstore-server/src/endpoints/resumable.rs +++ b/objectstore-server/src/endpoints/resumable.rs @@ -8,8 +8,8 @@ //! |---|---|---| //! | Create | `POST /objects/{usecase}/{scopes}/?upload_type=resumable` | `200` + `{"key","session","granularity"}` | //! | Create | `PUT /objects/{usecase}/{scopes}/{key}?upload_type=resumable` | `200` + `{"key","session","granularity"}` | -//! | Chunk | `PUT …/{key}?session=` with `Upload-Offset: ` | `204` + progress headers, or `201` + `{"key"}` | -//! | Offset query | `PUT …/{key}?session=` with `Upload-Offset: *` | `204` + progress headers, or `201` + `{"key"}` | +//! | Chunk | `PUT …/{key}?session=` with `Upload-Offset: ` | `204` + `Upload-Offset`, or `201` + `{"key"}` | +//! | Offset query | `PUT …/{key}?session=` with `Upload-Offset: *` | `204` + `Upload-Offset`, or `201` + `{"key"}` | //! | Cancel | `DELETE …/{key}?session=` | `204` | #![expect( @@ -28,8 +28,8 @@ use objectstore_service::id::{ObjectContext, ObjectId}; use objectstore_service::stream::ClientStream; use objectstore_types::metadata::Metadata; use objectstore_types::resumable::{ - CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_GRANULARITY, HEADER_UPLOAD_OFFSET, - UploadOffset, UploadProgress, + CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_OFFSET, UploadOffset, + UploadProgress, }; use objectstore_types::time::Timestamp; @@ -135,7 +135,7 @@ async fn create_session_for_id( /// /// Both answer `204 No Content` with the authoritative offset while bytes remain, and /// `201 Created` with the key once the upload is complete, the session is terminal, and the object -/// is available through the normal object endpoints. Both include `Upload-Granularity`. +/// is available through the normal object endpoints. /// The acknowledged offset may be lower than the submitted chunk's end, so clients should continue /// from this response (or from a later explicit offset query), never from local byte accounting alone. pub(super) async fn continue_session( @@ -147,8 +147,6 @@ pub(super) async fn continue_session( MeteredBody(body): MeteredBody, ) -> ApiResult { let key = id.key().to_owned(); - let granularity = service.upload_granularity(); - let progress = match offset { UploadOffset::At(offset) => { let content_length = content_length @@ -172,7 +170,7 @@ pub(super) async fn continue_session( } }; - progress_response(progress, key, granularity) + progress_response(progress, key) } /// Cancels a session, discarding whatever was uploaded. @@ -186,12 +184,8 @@ pub(super) async fn cancel_session( } /// Turns an [`UploadProgress`] outcome into the response shared by chunks and offset queries. -fn progress_response( - progress: ApiResult, - key: String, - granularity: u64, -) -> ApiResult { - let mut response = match progress { +fn progress_response(progress: ApiResult, key: String) -> ApiResult { + let response = match progress { Ok(UploadProgress::Incomplete { offset }) => ( StatusCode::NO_CONTENT, [(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset))], @@ -213,10 +207,6 @@ fn progress_response( }, Err(error) => return Err(error), }; - response.headers_mut().insert( - HEADER_UPLOAD_GRANULARITY, - http::HeaderValue::from(granularity), - ); Ok(response) } @@ -224,51 +214,38 @@ fn progress_response( mod tests { use super::*; - /// Reads a response's status, progress headers, and body. - async fn parts_of(response: Response) -> (StatusCode, Option, Option, String) { + /// Reads a response's status, offset header, and body. + async fn parts_of(response: Response) -> (StatusCode, Option, String) { let status = response.status(); let offset = response .headers() .get(HEADER_UPLOAD_OFFSET) .map(|v| v.to_str().unwrap().to_owned()); - let granularity = response - .headers() - .get(HEADER_UPLOAD_GRANULARITY) - .map(|v| v.to_str().unwrap().to_owned()); - let body = axum::body::to_bytes(response.into_body(), usize::MAX) .await .unwrap(); - ( - status, - offset, - granularity, - String::from_utf8(body.to_vec()).unwrap(), - ) + (status, offset, String::from_utf8(body.to_vec()).unwrap()) } #[tokio::test] async fn incomplete_progress_answers_no_content_with_the_offset() { let progress = Ok(UploadProgress::Incomplete { offset: 262_144 }); - let response = progress_response(progress, "my-key".into(), 262_144).unwrap(); + let response = progress_response(progress, "my-key".into()).unwrap(); - let (status, offset, granularity, body) = parts_of(response).await; + let (status, offset, body) = parts_of(response).await; assert_eq!(status, StatusCode::NO_CONTENT); assert_eq!(offset.as_deref(), Some("262144")); - assert_eq!(granularity.as_deref(), Some("262144")); assert!(body.is_empty(), "204 must not carry a body: {body:?}"); } #[tokio::test] async fn commit_answers_created_with_the_key() { - let response = - progress_response(Ok(UploadProgress::Complete), "my-key".into(), 262_144).unwrap(); + let response = progress_response(Ok(UploadProgress::Complete), "my-key".into()).unwrap(); - let (status, offset, granularity, body) = parts_of(response).await; + let (status, offset, body) = parts_of(response).await; assert_eq!(status, StatusCode::CREATED); assert_eq!(offset, None, "a commit reports no offset"); - assert_eq!(granularity.as_deref(), Some("262144")); assert_eq!(body, r#"{"key":"my-key"}"#); } @@ -276,23 +253,22 @@ mod tests { async fn offset_mismatch_answers_conflict_with_the_authoritative_offset() { let mismatch = ErrorKind::UploadOffsetMismatch { offset: 786_432 }.into(); let response = - progress_response(Err(ApiError::Service(mismatch)), "my-key".into(), 262_144).unwrap(); + progress_response(Err(ApiError::Service(mismatch)), "my-key".into()).unwrap(); - let (status, offset, granularity, body) = parts_of(response).await; + let (status, offset, body) = parts_of(response).await; assert_eq!(status, StatusCode::CONFLICT); assert_eq!( offset.as_deref(), Some("786432"), "the client resynchronizes from this header" ); - assert_eq!(granularity.as_deref(), Some("262144")); assert!(body.contains("786432"), "{body:?}"); } #[tokio::test] async fn other_errors_propagate_unchanged() { let gone = ApiError::Service(ErrorKind::UploadSessionGone.into()); - let error = progress_response(Err(gone), "my-key".into(), 262_144).unwrap_err(); + let error = progress_response(Err(gone), "my-key".into()).unwrap_err(); assert_eq!(error.status(), StatusCode::GONE); let oversized = @@ -301,7 +277,7 @@ mod tests { content_length: 4, upload_length: 10, })); - let error = progress_response(Err(oversized), "my-key".into(), 262_144).unwrap_err(); + let error = progress_response(Err(oversized), "my-key".into()).unwrap_err(); assert_eq!(error.status(), StatusCode::BAD_REQUEST); let too_small = ApiError::Service( @@ -311,7 +287,7 @@ mod tests { } .into(), ); - let error = progress_response(Err(too_small), "my-key".into(), 262_144).unwrap_err(); + let error = progress_response(Err(too_small), "my-key".into()).unwrap_err(); assert_eq!(error.status(), StatusCode::BAD_REQUEST); } } diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index 461356be..d4fc6ecc 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -8,7 +8,7 @@ use anyhow::Result; use objectstore_server::config::{AuthZ, Config, EncryptionConfig, Service}; use objectstore_test::server::TestServer; use objectstore_types::resumable::{ - CreateSessionResponse, HEADER_UPLOAD_GRANULARITY, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, + CreateSessionResponse, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, }; use reqwest::StatusCode; @@ -192,7 +192,6 @@ async fn test_resumable_upload() -> Result<()> { .await?; assert_eq!(response.status(), StatusCode::NO_CONTENT); assert_eq!(response.headers()[HEADER_UPLOAD_OFFSET], "3"); - assert_eq!(response.headers()[HEADER_UPLOAD_GRANULARITY], "0"); // Rejected: outdated offset. let response = client @@ -221,7 +220,6 @@ async fn test_resumable_upload() -> Result<()> { .await?; assert_eq!(response.status(), StatusCode::NO_CONTENT); assert_eq!(response.headers()[HEADER_UPLOAD_OFFSET], "3"); - assert_eq!(response.headers()[HEADER_UPLOAD_GRANULARITY], "0"); let response = client .put(server.url(&session_path)) .header(HEADER_UPLOAD_OFFSET, "3") diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 1fc52d2e..aa3e631e 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -334,13 +334,13 @@ Not all backends support resumable uploads. A backend returns no session when it particular upload; this is a routine outcome rather than an error. Acceptance can depend on the declared size, the metadata, or whether resuming is possible in principle. -The offset returned by every chunk response is authoritative. Each upload reports its granularity -in bytes, or zero when it imposes none. For a positive granularity, a backend may accept -the full request body but persist only the largest aligned prefix, so a client must not advance by -the submitted `Content-Length` on its own. Non-empty, non-final chunks shorter than one granularity -unit are rejected; the final chunk is exempt. Clients should use -granularity-sized multiples and continue from the preceding response's offset, or perform an -explicit offset query after an ambiguous failure. +The offset returned by every chunk response is authoritative. Session creation reports each +upload's granularity in bytes, or zero when it imposes none. For a positive granularity, a backend +may accept the full request body but persist only the largest aligned prefix, so a client must not +advance by the submitted `Content-Length` on its own. Non-empty, non-final chunks shorter than one +granularity unit are rejected; the final chunk is exempt. Clients should use granularity-sized +multiples and continue from the preceding response's offset, or perform an explicit offset query +after an ambiguous failure. ## Multipart Uploads diff --git a/objectstore-types/src/resumable.rs b/objectstore-types/src/resumable.rs index d71a55ee..c892cef7 100644 --- a/objectstore-types/src/resumable.rs +++ b/objectstore-types/src/resumable.rs @@ -3,8 +3,7 @@ //! A resumable upload writes one object across multiple requests. The client first creates a //! session, declaring the object's complete size with [`HEADER_UPLOAD_LENGTH`]. The server returns //! a [`CreateSessionResponse`] containing an opaque [`SessionToken`] that identifies the upload. -//! The response also reports the upload granularity. Subsequent progress responses -//! repeat that value in [`HEADER_UPLOAD_GRANULARITY`], so reconstructed clients can learn it. +//! The response also reports the upload granularity. A reconstructed client does not know it. //! //! The client then sends chunks with [`HEADER_UPLOAD_OFFSET`] set to the byte position at which //! each chunk starts. If an upload is interrupted, the client can send the wildcard offset @@ -34,13 +33,6 @@ pub const HEADER_UPLOAD_LENGTH: &str = "upload-length"; /// persisted. See [`UploadOffset`]. pub const HEADER_UPLOAD_OFFSET: &str = "upload-offset"; -/// Response header declaring the upload granularity, in bytes. -/// -/// A value of zero means the upload has no granularity. For a positive value, -/// non-final chunks shorter than one granularity unit are rejected, and a backend may persist -/// only a multiple of this value from a larger non-final chunk. -pub const HEADER_UPLOAD_GRANULARITY: &str = "upload-granularity"; - /// The wildcard [`HEADER_UPLOAD_OFFSET`] value that queries the server's offset. const OFFSET_WILDCARD: &str = "*"; From b13f886b6433b5aed46c36b70fef39c0f16013c0 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:01:08 +0200 Subject: [PATCH 07/18] wip --- clients/rust/src/resumable.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/clients/rust/src/resumable.rs b/clients/rust/src/resumable.rs index 2e624f41..831ea596 100644 --- a/clients/rust/src/resumable.rs +++ b/clients/rust/src/resumable.rs @@ -106,10 +106,10 @@ impl ResumableUpload { &self.token } - /// Returns this upload's granularity, in bytes. + /// Returns this upload's granularity in bytes, if known. /// - /// This is `None` for a reconstructed handle; only session creation supplies the value. - /// Zero means the upload has no granularity. + /// The granularity is the persistence unit for non-final chunks: chunks shorter than one + /// unit are rejected, and larger chunks may persist only an aligned prefix. pub fn granularity(&self) -> Option { self.granularity } From f6d1db2458896250f524d687abf62c64763422f4 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:03:38 +0200 Subject: [PATCH 08/18] improve --- objectstore-server/src/endpoints/resumable.rs | 29 ++++++++++++------- objectstore-service/src/backend/common.rs | 3 +- 2 files changed, 19 insertions(+), 13 deletions(-) diff --git a/objectstore-server/src/endpoints/resumable.rs b/objectstore-server/src/endpoints/resumable.rs index 8c527233..2b1b2f81 100644 --- a/objectstore-server/src/endpoints/resumable.rs +++ b/objectstore-server/src/endpoints/resumable.rs @@ -147,6 +147,7 @@ pub(super) async fn continue_session( MeteredBody(body): MeteredBody, ) -> ApiResult { let key = id.key().to_owned(); + let progress = match offset { UploadOffset::At(offset) => { let content_length = content_length @@ -185,15 +186,8 @@ pub(super) async fn cancel_session( /// Turns an [`UploadProgress`] outcome into the response shared by chunks and offset queries. fn progress_response(progress: ApiResult, key: String) -> ApiResult { - let response = match progress { - Ok(UploadProgress::Incomplete { offset }) => ( - StatusCode::NO_CONTENT, - [(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset))], - ) - .into_response(), - Ok(UploadProgress::Complete) => { - (StatusCode::CREATED, Json(CompleteUploadResponse { key })).into_response() - } + let progress = match progress { + Ok(progress) => progress, Err(ApiError::Service(error)) => match error.kind() { ErrorKind::UploadOffsetMismatch { offset } => { let error = ApiError::Service(error); @@ -201,12 +195,24 @@ fn progress_response(progress: ApiResult, key: String) -> ApiRes response .headers_mut() .insert(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset)); - response + return Ok(response); } _ => return Err(ApiError::Service(error)), }, Err(error) => return Err(error), }; + + let response = match progress { + UploadProgress::Incomplete { offset } => ( + StatusCode::NO_CONTENT, + [(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset))], + ) + .into_response(), + UploadProgress::Complete => { + (StatusCode::CREATED, Json(CompleteUploadResponse { key })).into_response() + } + }; + Ok(response) } @@ -214,13 +220,14 @@ fn progress_response(progress: ApiResult, key: String) -> ApiRes mod tests { use super::*; - /// Reads a response's status, offset header, and body. + /// Reads a response's status, `Upload-Offset` header, and body. async fn parts_of(response: Response) -> (StatusCode, Option, String) { let status = response.status(); let offset = response .headers() .get(HEADER_UPLOAD_OFFSET) .map(|v| v.to_str().unwrap().to_owned()); + let body = axum::body::to_bytes(response.into_body(), usize::MAX) .await .unwrap(); diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index ea0a2221..e1c31d61 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -194,8 +194,7 @@ pub trait Backend: fmt::Debug + Send + Sync + 'static { /// /// A value of zero means these uploads have no granularity. A positive value means that a /// non-final chunk can persist only a multiple of this value. Implementations must reject a - /// non-empty, non-final chunk shorter than one unit with [`ErrorKind::ChunkTooSmall`]. Final - /// chunks and zero-byte offset probes are exempt. + /// non-empty, non-final chunk shorter than one unit with [`ErrorKind::ChunkTooSmall`]. fn upload_granularity(&self) -> u64 { 0 } From 37ef6abbb02e40cb4ca8872c803fdc175c111e1a Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:04:43 +0200 Subject: [PATCH 09/18] improve --- clients/rust/src/resumable.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/clients/rust/src/resumable.rs b/clients/rust/src/resumable.rs index 831ea596..b268b9db 100644 --- a/clients/rust/src/resumable.rs +++ b/clients/rust/src/resumable.rs @@ -366,8 +366,7 @@ impl PutChunkBuilder { /// Returns [`Error::ResumableUploadUnavailable`] when the session expired, was canceled, or /// could not be found. The upload must be restarted with a new session in that case. /// Returns [`Error::ChunkTooSmall`] before sending the request when this is known to be a - /// non-final chunk shorter than the upload granularity. Reconstructed handles cannot perform - /// this check because they do not know the upload length or granularity. + /// non-final chunk shorter than the upload granularity. /// /// ```rust,ignore /// let offset = match upload.put(offset, chunk).send().await { From 50097aadcbbf73f7c825286c097a0da327f2357b Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:11:42 +0200 Subject: [PATCH 10/18] ref(resumable): Return upload granularity with created session --- objectstore-server/src/auth/service.rs | 11 ++---- objectstore-server/src/endpoints/resumable.rs | 6 +-- objectstore-service/src/service.rs | 39 +++++++++++++------ 3 files changed, 35 insertions(+), 21 deletions(-) diff --git a/objectstore-server/src/auth/service.rs b/objectstore-server/src/auth/service.rs index b16a9ffc..95e03a1a 100644 --- a/objectstore-server/src/auth/service.rs +++ b/objectstore-server/src/auth/service.rs @@ -4,7 +4,9 @@ use objectstore_service::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; -use objectstore_service::service::{DeleteResponse, GetResponse, InsertResponse, MetadataResponse}; +use objectstore_service::service::{ + CreatedUploadSession, DeleteResponse, GetResponse, InsertResponse, MetadataResponse, +}; use objectstore_service::{ClientStream, StorageService}; use objectstore_types::auth::Permission; @@ -218,18 +220,13 @@ impl AuthAwareService { // --- Resumable upload operations --- - /// Returns the upload granularity currently reported for new sessions, in bytes. - pub fn upload_granularity(&self) -> u64 { - self.service.upload_granularity() - } - /// Auth-aware wrapper around [`StorageService::create_upload_session`]. pub async fn create_upload_session( &self, id: ObjectId, metadata: Metadata, total_length: u64, - ) -> ApiResult> { + ) -> ApiResult> { self.check_permission(Permission::ObjectWrite, id.context())?; Ok(self .service diff --git a/objectstore-server/src/endpoints/resumable.rs b/objectstore-server/src/endpoints/resumable.rs index 2b1b2f81..1b72edd7 100644 --- a/objectstore-server/src/endpoints/resumable.rs +++ b/objectstore-server/src/endpoints/resumable.rs @@ -110,15 +110,15 @@ async fn create_session_for_id( .usecases .validate(&id.context().usecase, &metadata)?; - let session = service + let created = service .create_upload_session(id.clone(), metadata, total_length) .await? .ok_or_else(|| ServiceError::from(ErrorKind::Unsupported))?; let body = Json(CreateSessionResponse { key: id.key().to_owned(), - session, - granularity: service.upload_granularity(), + session: created.session, + granularity: created.granularity, }); Ok((StatusCode::OK, body).into_response()) } diff --git a/objectstore-service/src/service.rs b/objectstore-service/src/service.rs index 8a16caee..c0856ddc 100644 --- a/objectstore-service/src/service.rs +++ b/objectstore-service/src/service.rs @@ -43,6 +43,15 @@ pub type InsertResponse = ObjectId; /// Service response for [`StorageService::delete_object`]. pub type DeleteResponse = (); +/// A newly opened resumable upload session and its upload granularity. +#[derive(Debug)] +pub struct CreatedUploadSession { + /// The encrypted token used to continue the upload. + pub session: EncryptedSessionToken, + /// The upload granularity in bytes, or zero when there is none. + pub granularity: u64, +} + /// Default concurrency limit for [`StorageService`]. /// /// This value is used when no explicit limiter is set via @@ -465,11 +474,6 @@ impl StorageService { // --- Resumable upload operations --- - /// Returns the upload granularity currently reported for new sessions, in bytes. - pub fn upload_granularity(&self) -> u64 { - self.inner.upload_granularity() - } - /// Opens a resumable upload session for an object of `total_length` bytes. /// /// Returns `Ok(None)` for zero-length objects or when the backend declines resumable uploads @@ -479,7 +483,7 @@ impl StorageService { id: ObjectId, metadata: Metadata, total_length: u64, - ) -> Result> { + ) -> Result> { let Some(total_length) = NonZeroU64::new(total_length) else { return Ok(None); }; @@ -492,12 +496,16 @@ impl StorageService { .await?; session .map(|backend_token| { - cipher + let session = cipher .encrypt(&SessionToken { object_id: id, backend_token, }) - .map(EncryptedSessionToken::new) + .map(EncryptedSessionToken::new)?; + Ok(CreatedUploadSession { + session, + granularity: inner.upload_granularity(), + }) }) .transpose() }) @@ -605,6 +613,10 @@ mod tests { #[async_trait::async_trait] impl Hooks for ResumableTokenHooks { + fn upload_granularity(&self, _inner: &InMemoryBackend) -> u64 { + 256 * 1024 + } + async fn create_upload_session( &self, _inner: &InMemoryBackend, @@ -1400,11 +1412,13 @@ mod tests { async fn resumable_round_trip() { let service = make_service(); let id = ObjectId::new(make_context(), "resumable".into()); - let token = service + let created = service .create_upload_session(id.clone(), Metadata::default(), 3) .await .unwrap() .unwrap(); + assert_eq!(created.granularity, 0); + let token = created.session; assert_eq!( service .put_chunk(id.clone(), token.clone(), 0, 1, stream::single("a")) @@ -1478,10 +1492,12 @@ mod tests { ); let id = ObjectId::new(make_context(), "resumable".into()); - let token = service + let created = service .create_upload_session(id.clone(), Metadata::default(), 4) .await? .expect("test backend supports resumable uploads"); + assert_eq!(created.granularity, 256 * 1024); + let token = created.session; assert_ne!(token.as_bytes(), b"backend token"); assert!(matches!( service @@ -1517,7 +1533,8 @@ mod tests { let encrypted = service .create_upload_session(id.clone(), Metadata::default(), 4) .await? - .expect("test backend supports resumable uploads"); + .expect("test backend supports resumable uploads") + .session; assert_ne!(encrypted.as_bytes(), b"backend token"); service.upload_offset(id, encrypted).await?; assert_eq!( From 5173e57f22e80892524d1ca47015d0ead32acb73 Mon Sep 17 00:00:00 2001 From: Lorenzo Cian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:13:56 +0200 Subject: [PATCH 11/18] Apply suggestion from @lcian --- objectstore-server/tests/resumable.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index d4fc6ecc..bbfd5918 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -87,7 +87,7 @@ async fn create_session_path( .await?; assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; - assert_eq!(created.granularity, 0); + assert!(created.granularity.is_some()); Ok(format!( "{OBJECT_PATH}?session={}", created.session.to_base64url() From 3023b5f83d72e69aeb1626466fafbd3b95a57ed2 Mon Sep 17 00:00:00 2001 From: Lorenzo Cian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:14:09 +0200 Subject: [PATCH 12/18] Apply suggestion from @lcian --- objectstore-server/tests/resumable.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index bbfd5918..d14aaca5 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -173,7 +173,7 @@ async fn test_resumable_upload() -> Result<()> { .header(reqwest::header::CONTENT_TYPE, "text/plain") .send() .await?; - assert_eq!(response.status(), StatusCode::OK); + assert!(created.granularity.is_some()); let created: CreateSessionResponse = response.json().await?; assert_eq!(created.key, "my-key"); assert_eq!(created.granularity, 0); From 9e59eb5fe236924e7eef0b126434e9b8b7fb85a3 Mon Sep 17 00:00:00 2001 From: Lorenzo Cian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:14:44 +0200 Subject: [PATCH 13/18] Apply suggestion from @lcian --- objectstore-server/tests/resumable.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index d14aaca5..bbfd5918 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -173,7 +173,7 @@ async fn test_resumable_upload() -> Result<()> { .header(reqwest::header::CONTENT_TYPE, "text/plain") .send() .await?; - assert!(created.granularity.is_some()); + assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; assert_eq!(created.key, "my-key"); assert_eq!(created.granularity, 0); From 00902dca20c4d21492ae9e95d93ce72d3976dd47 Mon Sep 17 00:00:00 2001 From: Lorenzo Cian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:15:00 +0200 Subject: [PATCH 14/18] Apply suggestion from @lcian --- objectstore-server/tests/resumable.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index bbfd5918..b5380436 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -176,7 +176,7 @@ async fn test_resumable_upload() -> Result<()> { assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; assert_eq!(created.key, "my-key"); - assert_eq!(created.granularity, 0); + assert!(created.granularity.is_some()); let session_path = format!("{object}?session={}", created.session.to_base64url()); // The object doesn't exist yet. From 59a92eb24c77b58490dcbe966b2f14681a7a6ec0 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:20:41 +0200 Subject: [PATCH 15/18] test(resumable): Assert required upload granularity --- objectstore-server/tests/resumable.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/objectstore-server/tests/resumable.rs b/objectstore-server/tests/resumable.rs index b5380436..d4fc6ecc 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -87,7 +87,7 @@ async fn create_session_path( .await?; assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; - assert!(created.granularity.is_some()); + assert_eq!(created.granularity, 0); Ok(format!( "{OBJECT_PATH}?session={}", created.session.to_base64url() @@ -176,7 +176,7 @@ async fn test_resumable_upload() -> Result<()> { assert_eq!(response.status(), StatusCode::OK); let created: CreateSessionResponse = response.json().await?; assert_eq!(created.key, "my-key"); - assert!(created.granularity.is_some()); + assert_eq!(created.granularity, 0); let session_path = format!("{object}?session={}", created.session.to_base64url()); // The object doesn't exist yet. From 5a78d7bc50191919cf927fe616235f4d6c564270 Mon Sep 17 00:00:00 2001 From: Lorenzo Cian <17258265+lcian@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:22:18 +0200 Subject: [PATCH 16/18] Apply suggestion from @lcian --- clients/rust/tests/e2e.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/clients/rust/tests/e2e.rs b/clients/rust/tests/e2e.rs index c9234b1b..b5b3613e 100644 --- a/clients/rust/tests/e2e.rs +++ b/clients/rust/tests/e2e.rs @@ -1053,7 +1053,7 @@ async fn test_resumable_upload() { .unwrap() .unwrap(); assert_eq!(upload.key(), "resumable-client"); - assert_eq!(upload.granularity(), Some(0)); + assert!(upload.granularity().is_some()); assert_eq!( upload.progress().send().await.unwrap(), UploadProgress::Incomplete { offset: 0 } From febd1523a43cb2ea000351c8438478b8797b3944 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 29 Sep 2026 14:36:33 +0200 Subject: [PATCH 17/18] fmt --- objectstore-service/src/backend/tiered.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 1aa55d04..77742646 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -472,7 +472,7 @@ impl Backend for TieredStorage { content_length, upload_length: session.total_length, })?; - + let granularity = self.upload_granularity(); if content_length > 0 && content_length < granularity && end != session.total_length { return Err(ErrorKind::ChunkTooSmall { @@ -1320,7 +1320,7 @@ mod tests { let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1]; let token = resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await; - + let error = storage .put_chunk(&id, &token, 0, 1, stream::single("a")) .await From bb134741e75792704593abb6610bff7cb82ed8d6 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:36:46 +0200 Subject: [PATCH 18/18] improve --- clients/rust/README.md | 8 +++++--- objectstore-server/src/auth/service.rs | 4 ++-- objectstore-service/src/service.rs | 6 +++--- objectstore-types/src/resumable.rs | 3 ++- 4 files changed, 12 insertions(+), 9 deletions(-) diff --git a/clients/rust/README.md b/clients/rust/README.md index 18f40bcc..51509de7 100644 --- a/clients/rust/README.md +++ b/clients/rust/README.md @@ -167,9 +167,11 @@ If the request fails midway, it will be possible to resume it from the persisted **Important:** resumable uploads do not automatically compress chunk contents. The `compression` setting only records how the object is encoded; the caller must compress the payload accordingly. The object length and all offsets refer to the bytes after compression. -When manually slicing non-final chunks, use multiples of a positive upload granularity -(`upload.granularity()` on a newly created handle). Reconstructed handles do not know this value. -Shorter chunks are rejected, while a larger unaligned chunk may persist only its aligned prefix. + +When manually slicing non-final chunks, use sizes that are multiples of `upload.granularity()`. +All chunks except the last must be at least as large as the granularity. +If a non-final chunk is larger than the granularity but its size is not a multiple of it, then the +server may persist only its aligned prefix. ```rust,no_run #![cfg(feature = "resumable-upload-api")] diff --git a/objectstore-server/src/auth/service.rs b/objectstore-server/src/auth/service.rs index 95e03a1a..8f7b6513 100644 --- a/objectstore-server/src/auth/service.rs +++ b/objectstore-server/src/auth/service.rs @@ -5,7 +5,7 @@ use objectstore_service::multipart::{ ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; use objectstore_service::service::{ - CreatedUploadSession, DeleteResponse, GetResponse, InsertResponse, MetadataResponse, + CreateUploadSessionResponse, DeleteResponse, GetResponse, InsertResponse, MetadataResponse, }; use objectstore_service::{ClientStream, StorageService}; @@ -226,7 +226,7 @@ impl AuthAwareService { id: ObjectId, metadata: Metadata, total_length: u64, - ) -> ApiResult> { + ) -> ApiResult> { self.check_permission(Permission::ObjectWrite, id.context())?; Ok(self .service diff --git a/objectstore-service/src/service.rs b/objectstore-service/src/service.rs index c0856ddc..6082dfbe 100644 --- a/objectstore-service/src/service.rs +++ b/objectstore-service/src/service.rs @@ -45,7 +45,7 @@ pub type DeleteResponse = (); /// A newly opened resumable upload session and its upload granularity. #[derive(Debug)] -pub struct CreatedUploadSession { +pub struct CreateUploadSessionResponse { /// The encrypted token used to continue the upload. pub session: EncryptedSessionToken, /// The upload granularity in bytes, or zero when there is none. @@ -483,7 +483,7 @@ impl StorageService { id: ObjectId, metadata: Metadata, total_length: u64, - ) -> Result> { + ) -> Result> { let Some(total_length) = NonZeroU64::new(total_length) else { return Ok(None); }; @@ -502,7 +502,7 @@ impl StorageService { backend_token, }) .map(EncryptedSessionToken::new)?; - Ok(CreatedUploadSession { + Ok(CreateUploadSessionResponse { session, granularity: inner.upload_granularity(), }) diff --git a/objectstore-types/src/resumable.rs b/objectstore-types/src/resumable.rs index c892cef7..f3e3e0ec 100644 --- a/objectstore-types/src/resumable.rs +++ b/objectstore-types/src/resumable.rs @@ -185,7 +185,8 @@ pub struct CreateSessionResponse { pub key: String, /// The opaque session token that identifies the session. pub session: SessionToken, - /// This upload's granularity in bytes. + /// This upload's granularity in bytes, or zero when it imposes none. + #[serde(default)] pub granularity: u64, }