diff --git a/clients/rust/README.md b/clients/rust/README.md index 0d814b17..51509de7 100644 --- a/clients/rust/README.md +++ b/clients/rust/README.md @@ -168,6 +168,11 @@ 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 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/clients/rust/src/error.rs b/clients/rust/src/error.rs index 18d71a4b..a7fd6a49 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 chunk is shorter than the upload granularity. + #[cfg(feature = "resumable-upload-api")] + #[error( + "non-final chunk {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..b268b9db 100644 --- a/clients/rust/src/resumable.rs +++ b/clients/rust/src/resumable.rs @@ -52,6 +52,8 @@ pub struct ResumableUpload { session: Session, key: ObjectKey, token: SessionToken, + total_length: Option, + granularity: Option, } impl Session { @@ -87,6 +89,8 @@ impl Session { session: self.clone(), key: key.into(), token, + total_length: None, + granularity: None, } } } @@ -102,6 +106,14 @@ impl ResumableUpload { &self.token } + /// Returns this upload's granularity in bytes, if known. + /// + /// 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 + } + /// Builds a request for the server's authoritative upload progress. pub fn progress(&self) -> UploadProgressBuilder { UploadProgressBuilder { @@ -150,6 +162,28 @@ impl ResumableUpload { } } + fn validate_chunk_length(&self, offset: u64, chunk_length: u64) -> crate::Result<()> { + 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. pub fn cancel(&self) -> CancelUploadBuilder { CancelUploadBuilder { @@ -267,9 +301,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: Some(response.granularity), + }; + Ok(Some(upload)) } } @@ -326,6 +365,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 +377,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)? diff --git a/clients/rust/tests/e2e.rs b/clients/rust/tests/e2e.rs index 71b96c88..b5b3613e 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!(upload.granularity().is_some()); 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(), None); 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..8f7b6513 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::{ + CreateUploadSessionResponse, DeleteResponse, GetResponse, InsertResponse, MetadataResponse, +}; use objectstore_service::{ClientStream, StorageService}; use objectstore_types::auth::Permission; @@ -224,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-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..0b1a25d0 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 @@ -63,6 +64,8 @@ //! 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 +76,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..1b72edd7 100644 --- a/objectstore-server/src/endpoints/resumable.rs +++ b/objectstore-server/src/endpoints/resumable.rs @@ -6,8 +6,8 @@ //! //! | 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"}` | +//! | 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` + `Upload-Offset`, or `201` + `{"key"}` | //! | Offset query | `PUT …/{key}?session=` with `Upload-Offset: *` | `204` + `Upload-Offset`, or `201` + `{"key"}` | //! | Cancel | `DELETE …/{key}?session=` | `204` | @@ -110,14 +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, + session: created.session, + granularity: created.granularity, }); Ok((StatusCode::OK, body).into_response()) } @@ -285,5 +286,15 @@ mod tests { })); let error = progress_response(Err(oversized), "my-key".into()).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()).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..d4fc6ecc 100644 --- a/objectstore-server/tests/resumable.rs +++ b/objectstore-server/tests/resumable.rs @@ -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. diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 0009e067..f195610f 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -332,10 +332,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. 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-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index 8087a914..e1c31d61 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -190,6 +190,15 @@ 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`]. + fn upload_granularity(&self) -> u64 { + 0 + } + /// Stores an object at the given path with the given metadata. async fn put_object( &self, @@ -296,7 +305,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 bdf7b5c1..77742646 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) } @@ -469,6 +473,15 @@ impl Backend for TieredStorage { 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 { let progress = self @@ -1308,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..6082dfbe 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 CreateUploadSessionResponse { + /// 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 @@ -474,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); }; @@ -487,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(CreateUploadSessionResponse { + session, + granularity: inner.upload_granularity(), + }) }) .transpose() }) @@ -600,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, @@ -1395,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")) @@ -1473,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 @@ -1512,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!( diff --git a/objectstore-types/src/resumable.rs b/objectstore-types/src/resumable.rs index d1b6ea1d..f3e3e0ec 100644 --- a/objectstore-types/src/resumable.rs +++ b/objectstore-types/src/resumable.rs @@ -3,6 +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. 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 @@ -184,6 +185,9 @@ pub struct CreateSessionResponse { pub key: String, /// The opaque session token that identifies the session. pub session: SessionToken, + /// This upload's granularity in bytes, or zero when it imposes none. + #[serde(default)] + pub granularity: u64, } /// Response from the request that completes the upload. @@ -205,11 +209,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(()) }