Skip to content
Open
3 changes: 3 additions & 0 deletions clients/rust/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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()` 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
#![cfg(feature = "resumable-upload-api")]
Expand Down
11 changes: 11 additions & 0 deletions clients/rust/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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")]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'm noticing that our error type isn't marked as non_exhaustive, yet we add or remove variants with feature flags, which seems disruptive. By adding a feature flag, existing code may stop to compile.

When this stabilizes, I think we should keep those un-flagged and revisit our idea of hiding the function behind a feature flag permanently. It seems that it's better to just have the code in always, and instead:

  • either use an extension trait on session to make opt-in to the low-level methods explicit
  • or simply keep calling them out as low-level in the docs and that's it.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point. The ultimate cause of this though is that this enum is not marked non_exhaustive so whenever we introduce an error kind it will be a breaking change.
Should we consider marking this as non_exhaustive or would you rather do so when we know all the APIs are relatively stable?

#[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.
Expand Down
49 changes: 46 additions & 3 deletions clients/rust/src/resumable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,8 @@ pub struct ResumableUpload {
session: Session,
key: ObjectKey,
token: SessionToken,
total_length: Option<u64>,
granularity: Option<u64>,
}

impl Session {
Expand Down Expand Up @@ -87,6 +89,8 @@ impl Session {
session: self.clone(),
key: key.into(),
token,
total_length: None,
granularity: None,
}
}
}
Expand All @@ -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<u64> {
self.granularity
}

/// Builds a request for the server's authoritative upload progress.
pub fn progress(&self) -> UploadProgressBuilder {
UploadProgressBuilder {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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))
}
}

Expand Down Expand Up @@ -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 {
Expand All @@ -336,6 +377,8 @@ impl PutChunkBuilder {
/// };
/// ```
pub async fn send(self) -> crate::Result<UploadProgress> {
self.upload
.validate_chunk_length(self.offset, self.length)?;
let response = self
.upload
.request(Method::PUT)?
Expand Down
3 changes: 3 additions & 0 deletions clients/rust/tests/e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand All @@ -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 }
Expand Down
6 changes: 4 additions & 2 deletions objectstore-server/src/auth/service.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -224,7 +226,7 @@ impl AuthAwareService {
id: ObjectId,
metadata: Metadata,
total_length: u64,
) -> ApiResult<Option<SessionToken>> {
) -> ApiResult<Option<CreatedUploadSession>> {
self.check_permission(Permission::ObjectWrite, id.context())?;
Ok(self
.service
Expand Down
3 changes: 2 additions & 1 deletion objectstore-server/src/endpoints/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
9 changes: 6 additions & 3 deletions objectstore-server/src/endpoints/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
//!
Expand All @@ -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 |
Expand Down
19 changes: 15 additions & 4 deletions objectstore-server/src/endpoints/resumable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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=<s>` with `Upload-Offset: <n>` | `204` + `Upload-Offset`, or `201` + `{"key"}` |
//! | Offset query | `PUT …/{key}?session=<s>` with `Upload-Offset: *` | `204` + `Upload-Offset`, or `201` + `{"key"}` |
//! | Cancel | `DELETE …/{key}?session=<s>` | `204` |
Expand Down Expand Up @@ -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())
}
Expand Down Expand Up @@ -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);
}
}
2 changes: 2 additions & 0 deletions objectstore-server/tests/resumable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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.
Expand Down
11 changes: 7 additions & 4 deletions objectstore-service/docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. 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

Expand Down
12 changes: 11 additions & 1 deletion objectstore-service/src/backend/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
4 changes: 4 additions & 0 deletions objectstore-service/src/backend/counting.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading