diff --git a/Cargo.lock b/Cargo.lock index a08275d6..c8a90870 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2809,6 +2809,7 @@ dependencies = [ "argh", "async-stream", "axum", + "base64", "bytes", "ed25519-dalek", "elegant-departure", diff --git a/objectstore-server/Cargo.toml b/objectstore-server/Cargo.toml index 3e3d446b..13d4da82 100644 --- a/objectstore-server/Cargo.toml +++ b/objectstore-server/Cargo.toml @@ -15,6 +15,7 @@ anyhow = { workspace = true } argh = { workspace = true } async-stream = { workspace = true } axum = { workspace = true, features = ["multipart"] } +base64 = { workspace = true } bytes = { workspace = true } ed25519-dalek = { workspace = true, features = ["pem"] } elegant-departure = { workspace = true, features = ["tokio"] } diff --git a/objectstore-server/docs/architecture.md b/objectstore-server/docs/architecture.md index e665fc96..6dfc4fc0 100644 --- a/objectstore-server/docs/architecture.md +++ b/objectstore-server/docs/architecture.md @@ -7,49 +7,16 @@ core storage operations. ## Endpoints -All object operations live under the `/v1/` prefix: - -| Method | Path | Description | -|----------|-------------------------------------------|------------------------------| -| `POST` | `/v1/objects/{usecase}/{scopes}/` | Insert with server-generated key | -| `GET` | `/v1/objects/{usecase}/{scopes}/{*key}` | Retrieve object | -| `HEAD` | `/v1/objects/{usecase}/{scopes}/{*key}` | Retrieve metadata only | -| `PUT` | `/v1/objects/{usecase}/{scopes}/{*key}` | Insert or overwrite with key | -| `DELETE` | `/v1/objects/{usecase}/{scopes}/{*key}` | Delete object | -| `POST` | `/v1/objects:batch/{usecase}/{scopes}/` | Batch operations (multipart) | - -### Multipart Upload Endpoints - -| Method | Path | Description | -|-----------|--------------------------------------------------------------|--------------------------------------| -| `POST` | `/v1/objects:multipart/{usecase}/{scopes}/` | Initiate upload (server-generated key) | -| `PUT` | `/v1/objects:multipart/{usecase}/{scopes}/{*key}` | Initiate upload (user-provided key) | -| `PUT` | `/v1/objects:multipart:parts/{usecase}/{scopes}/{*key}` | Upload a part (`uploadId`, `partNumber` query params) | -| `GET` | `/v1/objects:multipart:parts/{usecase}/{scopes}/{*key}` | List uploaded parts (`uploadId` query param) | -| `POST` | `/v1/objects:multipart:complete/{usecase}/{scopes}/{*key}` | Complete upload (`uploadId` query param) | -| `DELETE` | `/v1/objects:multipart/{usecase}/{scopes}/{*key}` | Abort upload (`uploadId` query param) | - -The initiate POST endpoint accepts both trailing-slash and non-trailing-slash forms. - -The complete endpoint returns `200 OK` immediately, with a streaming body that -will contain the error (if any) as JSON. Whitespace is sent in the streaming body -to keep the connection open. -Clients must parse the body to determine the actual outcome, and not rely on the -status code. - -Scopes are encoded in the URL path using Matrix URI syntax: -`org=123;project=456`. An underscore (`_`) represents empty scopes. - -### Internal Endpoints - -Internal endpoints are exempt from authentication, rate limiting, and the web -concurrency limit so they remain available when the server is under load. - -| Method | Path | Description | -|--------|------|-------------| -| `GET` | `/health` | Liveness probe (always returns 200) | -| `GET` | `/ready` | Readiness probe (returns 503 when `/tmp/objectstore.down` exists, enabling graceful drain) | -| `GET` | `/keda` | Prometheus text-format gauges for KEDA autoscaling (see [KEDA Metrics](#keda-metrics)) | +All object operations live under the `/v1/` prefix. Objects are addressed by a +usecase, a set of scopes, and a key; scopes are encoded in the URL path using +Matrix URI syntax (`org=123;project=456`). + +Four families of routes exist: object operations, resumable uploads, multipart +uploads (being replaced by resumable uploads), and internal probes that stay +available when the server is under load. + +See the [`endpoints`] module for the routing table and the request and response +shape of every route. ## Request Flow diff --git a/objectstore-server/src/auth/service.rs b/objectstore-server/src/auth/service.rs index f4643bda..43c03701 100644 --- a/objectstore-server/src/auth/service.rs +++ b/objectstore-server/src/auth/service.rs @@ -3,6 +3,7 @@ use objectstore_service::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; +use objectstore_service::resumable::{SessionToken, UploadProgress}; use objectstore_service::service::{DeleteResponse, GetResponse, InsertResponse, MetadataResponse}; use objectstore_service::{ClientStream, StorageService}; @@ -186,4 +187,57 @@ impl AuthAwareService { .complete_multipart(id, upload_id, parts) .await?) } + + // --- Resumable upload operations --- + // + // Every operation requires `ObjectWrite`, including the two that do not obviously write: + // an offset query can commit an assembled object, and canceling a session discards an + // in-progress upload rather than deleting an object. So `DELETE ?session=` needs write + // permission where a plain `DELETE` on the same path needs delete permission. + + /// Auth-aware wrapper around [`StorageService::create_upload_session`]. + pub async fn create_upload_session( + &self, + id: ObjectId, + metadata: Metadata, + total_length: u64, + ) -> ApiResult { + self.check_permission(Permission::ObjectWrite, id.context())?; + Ok(self + .service + .create_upload_session(id, metadata, total_length) + .await?) + } + + /// Auth-aware wrapper around [`StorageService::put_chunk`]. + pub async fn put_chunk( + &self, + id: ObjectId, + session: SessionToken, + offset: u64, + content_length: u64, + body: ClientStream, + ) -> ApiResult { + self.check_permission(Permission::ObjectWrite, id.context())?; + Ok(self + .service + .put_chunk(id, session, offset, content_length, body) + .await?) + } + + /// Auth-aware wrapper around [`StorageService::upload_offset`]. + pub async fn upload_offset( + &self, + id: ObjectId, + session: SessionToken, + ) -> ApiResult { + self.check_permission(Permission::ObjectWrite, id.context())?; + Ok(self.service.upload_offset(id, session).await?) + } + + /// Auth-aware wrapper around [`StorageService::cancel_upload`]. + pub async fn cancel_upload(&self, id: ObjectId, session: SessionToken) -> ApiResult<()> { + self.check_permission(Permission::ObjectWrite, id.context())?; + Ok(self.service.cancel_upload(id, session).await?) + } } diff --git a/objectstore-server/src/endpoints/common.rs b/objectstore-server/src/endpoints/common.rs index 97468764..96fcde0e 100644 --- a/objectstore-server/src/endpoints/common.rs +++ b/objectstore-server/src/endpoints/common.rs @@ -96,6 +96,12 @@ impl ApiError { StatusCode::RANGE_NOT_SATISFIABLE } ApiError::Service(ServiceError::InvalidUploadId(_)) => StatusCode::BAD_REQUEST, + ApiError::Service(ServiceError::UnknownUploadSession) => StatusCode::BAD_REQUEST, + ApiError::Service(ServiceError::ChunkExceedsUploadLength { .. }) => { + StatusCode::BAD_REQUEST + } + ApiError::Service(ServiceError::UploadOffsetMismatch { .. }) => StatusCode::CONFLICT, + ApiError::Service(ServiceError::UploadSessionGone) => StatusCode::GONE, ApiError::Service(ServiceError::AtCapacity) => StatusCode::TOO_MANY_REQUESTS, ApiError::Service(ServiceError::NotImplemented) => StatusCode::NOT_IMPLEMENTED, ApiError::Service(_) => StatusCode::INTERNAL_SERVER_ERROR, diff --git a/objectstore-server/src/endpoints/mod.rs b/objectstore-server/src/endpoints/mod.rs index ab460b69..42a9755d 100644 --- a/objectstore-server/src/endpoints/mod.rs +++ b/objectstore-server/src/endpoints/mod.rs @@ -1,5 +1,116 @@ //! Contains all HTTP endpoint handlers. //! +//! This module documents the request and response shape of every route; see the [crate +//! documentation](crate) for the layers a request passes through before reaching a handler. +//! +//! Scopes are encoded in the URL path using Matrix URI syntax: `org=123;project=456`. An +//! underscore (`_`) represents empty scopes. +//! +//! # Object Endpoints +//! +//! All object operations live under the `/v1/` prefix: +//! +//! | Method | Path | Description | +//! |----------|-------------------------------------------|------------------------------| +//! | `POST` | `/v1/objects/{usecase}/{scopes}/` | Insert with server-generated key | +//! | `GET` | `/v1/objects/{usecase}/{scopes}/{*key}` | Retrieve object | +//! | `HEAD` | `/v1/objects/{usecase}/{scopes}/{*key}` | Retrieve metadata only | +//! | `PUT` | `/v1/objects/{usecase}/{scopes}/{*key}` | Insert or overwrite with key | +//! | `DELETE` | `/v1/objects/{usecase}/{scopes}/{*key}` | Delete object | +//! | `POST` | `/v1/objects:batch/{usecase}/{scopes}/` | Batch operations (multipart) | +//! +//! Object metadata travels in request and response headers; see +//! [`objectstore_types::metadata`] for the mapping. +//! +//! # Resumable Upload Endpoints +//! +//! A resumable upload transfers a single object across several requests. +//! The client opens a session, declaring the object's total size and metadata upfront, and +//! then sends the payload as a sequence of chunks at increasing byte offsets. +//! If a chunk fails, the client can ask the server which offset it holds and continue from there, +//! so an interrupted transfer resumes where it stopped instead of starting over. +//! Clients are still encouraged to send the whole payload in a single request, as that's the most +//! efficient and reliable approach. +//! The server knows the total size from the session, so it recognizes the chunk carrying the last +//! byte and commits the object itself. +//! +//! Resumable uploads use the object endpoints above, selected by a query parameter: +//! `upload_type=resumable` opens a session, and `session=` addresses it from then on. +//! Session creation returns the opaque backend token unchanged; subsequent requests encode it as +//! unpadded base64url in the `session` query parameter. +//! The object is named by the request path as usual, and [`objectstore_types::resumable`] +//! holds the protocol types. +//! +//! | Method | Path | Description | +//! |----------|------------------------------------------------------------|----------------------------------------------| +//! | `POST` | `/v1/objects/{usecase}/{scopes}/?upload_type=resumable` | Create session (server-generated key) | +//! | `PUT` | `/v1/objects/{usecase}/{scopes}/{*key}?upload_type=resumable` | Create session (user-provided key) | +//! | `PUT` | `/v1/objects/{usecase}/{scopes}/{*key}?session=` | Upload a chunk, or query the offset | +//! | `DELETE` | `/v1/objects/{usecase}/{scopes}/{*key}?session=` | Cancel upload, discarding what was sent | +//! +//! 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. +//! +//! 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 +//! the `*` wildcard submits an empty body and asks which offset the server holds. Both answer +//! `204 No Content` with the authoritative `Upload-Offset` while bytes remain, and +//! `201 Created` with `{"key"}` once the object is committed. +//! 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. +//! Every chunk requires `Content-Length`, even over HTTP/2, while creation and offset queries +//! must not carry a request body. +//! +//! An offset query can commit an object, so it requires write permission despite being read-shaped. +//! Termination likewise needs write rather than delete permission: it releases an in-progress upload, +//! not an object. +//! +//! | Status | Meaning | Client action | +//! |--------|---------|---------------| +//! | `400` | Malformed: unknown upload session, missing `Upload-Length`, nonempty offset query, or a chunk exceeding the declared length | Terminal | +//! | `409` | A chunk's offset does not match, with the authoritative offset in `Upload-Offset` | Resynchronize | +//! | `410` | The session expired or was canceled; nothing was retained | Start a new session | +//! | `501` | The configured backend does not implement resumable uploads | Fall back to a regular upload | +//! +//! # Multipart Upload Endpoints +//! +//! Multipart uploads are being replaced by [resumable +//! uploads](#resumable-upload-endpoints) and will be removed once all consumers have +//! migrated. See [`objectstore_types::multipart`] for the protocol types. +//! +//! | Method | Path | Description | +//! |-----------|--------------------------------------------------------------|--------------------------------------| +//! | `POST` | `/v1/objects:multipart/{usecase}/{scopes}/` | Initiate upload (server-generated key) | +//! | `PUT` | `/v1/objects:multipart/{usecase}/{scopes}/{*key}` | Initiate upload (user-provided key) | +//! | `PUT` | `/v1/objects:multipart:parts/{usecase}/{scopes}/{*key}` | Upload a part (`upload_id`, `part_number` query params) | +//! | `GET` | `/v1/objects:multipart:parts/{usecase}/{scopes}/{*key}` | List uploaded parts (`upload_id` query param) | +//! | `POST` | `/v1/objects:multipart:complete/{usecase}/{scopes}/{*key}` | Complete upload (`upload_id` query param) | +//! | `DELETE` | `/v1/objects:multipart/{usecase}/{scopes}/{*key}` | Abort upload (`upload_id` query param) | +//! +//! The initiate POST endpoint accepts both trailing-slash and non-trailing-slash forms. +//! +//! The complete endpoint returns `200 OK` immediately, with a streaming body that will +//! contain the error (if any) as JSON. Whitespace is sent in the streaming body to keep the +//! connection open. Clients must parse the body to determine the actual outcome, and not rely +//! on the status code. +//! +//! # Internal Endpoints +//! +//! Internal endpoints are exempt from authentication, rate limiting, and the web concurrency +//! limit so they remain available when the server is under load. [`is_internal_route`] +//! identifies them. +//! +//! | Method | Path | Description | +//! |--------|------|-------------| +//! | `GET` | `/health` | Liveness probe (always returns 200) | +//! | `GET` | `/ready` | Readiness probe (returns 503 when `/tmp/objectstore.down` exists, enabling graceful drain) | +//! | `GET` | `/keda` | Prometheus text-format gauges for KEDA autoscaling (see [KEDA Metrics](crate#keda-metrics)) | +//! +//! # Code Usage +//! //! Use [`routes`] to create a router with all endpoints. use axum::Router; @@ -14,6 +125,7 @@ mod multipart; mod objects; #[cfg(all(target_os = "linux", feature = "profiling"))] mod profiling; +mod resumable; /// Returns `true` for internal endpoints that are exempt from metrics and concurrency limits. pub fn is_internal_route(route: &str) -> bool { diff --git a/objectstore-server/src/endpoints/objects.rs b/objectstore-server/src/endpoints/objects.rs index 1aafcbac..b4fd0320 100644 --- a/objectstore-server/src/endpoints/objects.rs +++ b/objectstore-server/src/endpoints/objects.rs @@ -1,7 +1,8 @@ use std::fmt::Write as _; use axum::body::Body; -use axum::extract::State; +use axum::extract::{Request, State}; +use axum::handler::Handler; use axum::http::{HeaderMap, StatusCode}; use axum::response::{IntoResponse, Response}; use axum::routing; @@ -15,17 +16,18 @@ use serde::Serialize; use crate::auth::AuthAwareService; use crate::endpoints::common::{ApiError, ApiResult, insert_accept_ranges}; +use crate::endpoints::resumable::{self, ResumableTarget}; use crate::extractors::byte_range::OptionalByteRange; use crate::extractors::{Xt, body::MeteredBody}; use crate::state::ServiceState; pub fn router() -> Router { - let collection_routes = routing::post(objects_post); + let collection_routes = routing::post(dispatch_objects_post); let object_routes = routing::get(object_get) .head(object_head) - .put(object_put) + .put(dispatch_object_put) // TODO(ja): Implement PATCH (metadata update w/o body) - .delete(object_delete); + .delete(dispatch_object_delete); Router::new() .route("/objects/{usecase}/{scopes}", collection_routes.clone()) @@ -33,13 +35,60 @@ pub fn router() -> Router { .route("/objects/{usecase}/{scopes}/{*key}", object_routes) } +async fn dispatch_objects_post( + State(state): State, + target: Option, + request: Request, +) -> Response { + match target { + Some(ResumableTarget::NewSession) => resumable::create_session.call(request, state).await, + Some(ResumableTarget::ExistingSession) => { + ApiError::Client("`session` requires an object key; use PUT on the object path".into()) + .into_response() + } + None => create_object.call(request, state).await, + } +} + +async fn dispatch_object_put( + State(state): State, + target: Option, + request: Request, +) -> Response { + match target { + Some(ResumableTarget::NewSession) => { + resumable::create_session_for_key.call(request, state).await + } + Some(ResumableTarget::ExistingSession) => { + resumable::continue_session.call(request, state).await + } + None => insert_object.call(request, state).await, + } +} + +async fn dispatch_object_delete( + State(state): State, + target: Option, + request: Request, +) -> Response { + match target { + Some(ResumableTarget::ExistingSession) => { + resumable::cancel_session.call(request, state).await + } + Some(ResumableTarget::NewSession) => { + ApiError::Client("`upload_type` is not supported on DELETE".into()).into_response() + } + None => delete_object.call(request, state).await, + } +} + /// Response returned when inserting an object. #[derive(Debug, Serialize)] pub struct InsertObjectResponse { pub key: String, } -async fn objects_post( +async fn create_object( service: AuthAwareService, State(state): State, Xt(context): Xt, @@ -195,7 +244,7 @@ fn format_content_disposition(filename: &str) -> http::HeaderValue { http::HeaderValue::from_str(&result).expect("content disposition is a valid header value") } -async fn object_put( +async fn insert_object( service: AuthAwareService, State(state): State, Xt(id): Xt, @@ -223,7 +272,7 @@ async fn object_put( Ok((StatusCode::OK, response).into_response()) } -async fn object_delete( +async fn delete_object( service: AuthAwareService, Xt(id): Xt, ) -> ApiResult { diff --git a/objectstore-server/src/endpoints/resumable.rs b/objectstore-server/src/endpoints/resumable.rs new file mode 100644 index 00000000..a8d63c34 --- /dev/null +++ b/objectstore-server/src/endpoints/resumable.rs @@ -0,0 +1,481 @@ +//! Resumable upload endpoints. +//! +//! Resumable uploads are a variation of the regular object endpoints rather than a separate +//! resource. Every request addresses a standard object path with the session in the query string +//! (or `upload_type=resumable` for creation). +//! +//! | 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"}` | +//! | Cancel | `DELETE …/{key}?session=` | `204` | + +use axum::extract::{FromRequestParts, OptionalFromRequestParts, Query, State}; +use axum::http::{HeaderMap, StatusCode, request::Parts}; +use axum::response::{IntoResponse, Response}; +use axum::{Json, http}; +use base64::Engine as _; +use base64::engine::general_purpose::URL_SAFE_NO_PAD; +use futures_util::TryStreamExt; +use objectstore_service::error::Error as ServiceError; +use objectstore_service::id::{ObjectContext, ObjectId}; +use objectstore_service::resumable::{SessionToken, UploadOffset, UploadProgress}; +use objectstore_service::stream::ClientStream; +use objectstore_types::metadata::Metadata; +use objectstore_types::resumable::{ + CommitResponse, CreateSessionResponse, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET, +}; +use serde::Deserialize; + +use crate::auth::AuthAwareService; +use crate::endpoints::common::{ApiError, ApiResult}; +use crate::extractors::{Xt, body::MeteredBody}; +use crate::state::ServiceState; + +/// The `upload_type` query parameter. +#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub(super) enum UploadType { + /// Create a resumable upload session. + Resumable, +} + +/// The resumable protocol's query parameters, as seen on a regular object route. +#[derive(Debug, Deserialize)] +struct ResumableQuery { + /// Present on a session creation request. + upload_type: Option, + /// Present on a chunk write, offset query, or cancellation. + session: Option, +} + +/// Which resumable session a request on an object route targets. +#[derive(Debug)] +pub(super) enum ResumableTarget { + /// A new session to create for the object addressed by the request. + NewSession, + /// An existing session to continue or cancel. + ExistingSession, +} + +impl ResumableQuery { + /// Classifies a request that may create a session or act on one. + /// + /// # Errors + /// + /// Returns [`ApiError::Client`] if both parameters are present. + fn classify(self) -> ApiResult> { + match (self.upload_type, self.session) { + (Some(_), Some(_)) => Err(ApiError::Client( + "`upload_type` and `session` are mutually exclusive".into(), + )), + (Some(UploadType::Resumable), None) => Ok(Some(ResumableTarget::NewSession)), + (None, Some(_)) => Ok(Some(ResumableTarget::ExistingSession)), + (None, None) => Ok(None), + } + } +} + +/// A session token decoded by a continuation or cancellation handler. +#[derive(Debug)] +pub(super) struct Session(SessionToken); + +#[derive(Debug, Deserialize)] +struct SessionQuery { + session: String, +} + +/// Decodes the session token from its query-string representation. +fn decode_session_token(encoded: &str) -> ApiResult { + let bytes = URL_SAFE_NO_PAD + .decode(encoded) + .map_err(|error| ApiError::Client(error.to_string()))?; + + if URL_SAFE_NO_PAD.encode(&bytes) != encoded { + return Err(ApiError::Client( + "session token must use unpadded base64url encoding".into(), + )); + } + + String::from_utf8(bytes).map_err(|error| ApiError::Client(error.to_string())) +} + +impl FromRequestParts for Session +where + S: Send + Sync, +{ + type Rejection = ApiError; + + async fn from_request_parts(parts: &mut Parts, _state: &S) -> ApiResult { + let Query(SessionQuery { session }) = Query::::try_from_uri(&parts.uri) + .map_err(|error| ApiError::Client(error.to_string()))?; + Ok(Session(decode_session_token(&session)?)) + } +} + +impl OptionalFromRequestParts for ResumableTarget +where + S: Send + Sync, +{ + type Rejection = ApiError; + + async fn from_request_parts( + parts: &mut Parts, + _state: &S, + ) -> ApiResult> { + if parts.uri.query().is_none() { + return Ok(None); + } + + let Query(query) = Query::::try_from_uri(&parts.uri) + .map_err(|error| ApiError::Client(error.to_string()))?; + + query.classify() + } +} + +/// Reads the required [`HEADER_UPLOAD_LENGTH`] header. +fn upload_length(headers: &HeaderMap) -> ApiResult { + let value = headers + .get(HEADER_UPLOAD_LENGTH) + .ok_or_else(|| ApiError::Client(format!("{HEADER_UPLOAD_LENGTH} header is required")))?; + + value + .to_str() + .ok() + .filter(|v| v.bytes().all(|b| b.is_ascii_digit())) + .and_then(|v| v.parse().ok()) + .ok_or_else(|| ApiError::Client(format!("{HEADER_UPLOAD_LENGTH} must be a byte count"))) +} + +/// Reads the required [`HEADER_UPLOAD_OFFSET`] header. +fn upload_offset(headers: &HeaderMap) -> ApiResult { + let value = headers + .get(HEADER_UPLOAD_OFFSET) + .ok_or_else(|| ApiError::Client(format!("{HEADER_UPLOAD_OFFSET} header is required")))?; + + value + .to_str() + .map_err(|_| ApiError::Client(format!("{HEADER_UPLOAD_OFFSET} must be ASCII")))? + .parse() + .map_err(|e: objectstore_types::resumable::InvalidUploadOffset| { + ApiError::Client(e.to_string()) + }) +} + +/// Reads the required `Content-Length` header. +/// +/// Chunks declare their length so the server can forward only the prefix a backend accepts +/// without buffering the body to find out how long it is. +fn content_length(headers: &HeaderMap) -> ApiResult { + headers + .get(http::header::CONTENT_LENGTH) + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.parse::().ok()) + .ok_or_else(|| ApiError::Client("Content-Length header is required".into())) +} + +/// Confirms that a request neither declares nor streams a non-empty body. +async fn require_empty_body( + headers: &HeaderMap, + mut body: ClientStream, + request: &str, +) -> ApiResult<()> { + if headers.contains_key(http::header::CONTENT_LENGTH) && content_length(headers)? > 0 { + return Err(ApiError::Client(format!( + "{request} must be sent with an empty body" + ))); + } + + while let Some(chunk) = body.try_next().await.map_err(ServiceError::from)? { + if !chunk.is_empty() { + return Err(ApiError::Client(format!( + "{request} must be sent with an empty body" + ))); + } + } + + Ok(()) +} + +/// Creates a session with a server-generated object key. +pub(super) async fn create_session( + service: AuthAwareService, + State(state): State, + Xt(context): Xt, + headers: HeaderMap, + MeteredBody(body): MeteredBody, +) -> ApiResult { + create_session_for_id( + service, + state, + ObjectId::optional(context, None), + headers, + body, + ) + .await +} + +/// Creates a session for the object key in the request path. +pub(super) async fn create_session_for_key( + service: AuthAwareService, + State(state): State, + Xt(id): Xt, + headers: HeaderMap, + MeteredBody(body): MeteredBody, +) -> ApiResult { + create_session_for_id(service, state, id, headers, body).await +} + +/// Creates a session for the object at `id`. +/// +/// Answers `501 Not Implemented` when the backend refuses to create a session or doesn't implement +/// resumable uploads. +async fn create_session_for_id( + service: AuthAwareService, + state: ServiceState, + id: ObjectId, + headers: HeaderMap, + body: ClientStream, +) -> ApiResult { + let total_length = upload_length(&headers)?; + require_empty_body(&headers, body, "resumable session creation").await?; + let metadata = Metadata::from_insert_headers(&headers, "").map_err(ServiceError::from)?; + + state + .config + .usecases + .validate(&id.context().usecase, &metadata) + .map_err(|e| ApiError::Client(e.to_string()))?; + + let session = service + .create_upload_session(id.clone(), metadata, total_length) + .await?; + + let body = Json(CreateSessionResponse { + key: id.key().to_owned(), + session, + }); + Ok((StatusCode::OK, body).into_response()) +} + +/// Acts on an open session: writes a chunk, or reports the offset the server holds. +/// +/// [`HEADER_UPLOAD_OFFSET`] selects between the two. A concrete offset submits the request +/// body as the chunk starting there; the `*` wildcard submits nothing and asks where the +/// server stands, which also commits an object that was assembled but not yet committed. +/// Chunks require `Content-Length`, including over HTTP/2. Offset queries may omit it, but the +/// body stream is checked and any bytes are rejected as a malformed request. +/// +/// Both answer `204 No Content` with the authoritative offset while bytes remain, and +/// `201 Created` with the key once the object is committed. +pub(super) async fn continue_session( + service: AuthAwareService, + Xt(id): Xt, + Session(session): Session, + headers: HeaderMap, + MeteredBody(body): MeteredBody, +) -> ApiResult { + let offset = upload_offset(&headers)?; + let key = id.key().to_owned(); + + let progress = match offset { + UploadOffset::At(offset) => { + let content_length = content_length(&headers)?; + service + .put_chunk(id, session, offset, content_length, body) + .await + } + UploadOffset::Unknown => { + // The wildcard carries no payload. A body would be silently discarded, so + // reject it rather than let a client believe those bytes were written. + require_empty_body(&headers, body, "Upload-Offset: *").await?; + + service.upload_offset(id, session).await + } + }; + + progress_response(progress, key) +} + +/// Cancels a session, discarding whatever was uploaded. +pub(super) async fn cancel_session( + service: AuthAwareService, + Xt(id): Xt, + Session(session): Session, +) -> ApiResult { + service.cancel_upload(id, session).await?; + Ok(StatusCode::NO_CONTENT.into_response()) +} + +/// 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, + Err(error @ ApiError::Service(ServiceError::UploadOffsetMismatch { offset })) => { + let mut response = error.into_response(); + response + .headers_mut() + .insert(HEADER_UPLOAD_OFFSET, http::HeaderValue::from(offset)); + return Ok(response); + } + 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::Committed => { + (StatusCode::CREATED, Json(CommitResponse { key })).into_response() + } + }; + + Ok(response) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn query(upload_type: Option, session: Option<&str>) -> ResumableQuery { + ResumableQuery { + upload_type, + session: session.map(str::to_owned), + } + } + + #[test] + fn classify_recognizes_each_operation() { + assert!(matches!( + query(Some(UploadType::Resumable), None).classify(), + Ok(Some(ResumableTarget::NewSession)) + )); + assert!(matches!( + query(None, Some("token")).classify(), + Ok(Some(ResumableTarget::ExistingSession)) + )); + assert!(matches!(query(None, None).classify(), Ok(None))); + } + + #[test] + fn classify_rejects_both_parameters() { + let result = query(Some(UploadType::Resumable), Some("token")).classify(); + assert!(matches!(result, Err(ApiError::Client(_))), "{result:?}"); + } + + #[test] + fn session_token_decodes_from_unpadded_base64url() { + assert_eq!(decode_session_token("Li4vZXNjYXBl").unwrap(), "../escape"); + } + + #[test] + fn session_token_rejects_invalid_query_encodings() { + for invalid in ["%%%", "dG9rM24=", "_w"] { + assert!( + decode_session_token(invalid).is_err(), + "accepted {invalid:?}" + ); + } + } + + #[test] + fn upload_length_requires_a_byte_count() { + let mut headers = HeaderMap::new(); + assert!(upload_length(&headers).is_err(), "missing header"); + + for invalid in ["", "-1", "+1", "1.5", "abc", " 1"] { + headers.insert(HEADER_UPLOAD_LENGTH, invalid.parse().unwrap()); + assert!(upload_length(&headers).is_err(), "accepted {invalid:?}"); + } + + headers.insert(HEADER_UPLOAD_LENGTH, "1048576".parse().unwrap()); + assert_eq!(upload_length(&headers).unwrap(), 1_048_576); + } + + #[test] + fn upload_offset_parses_chunk_and_wildcard() { + let mut headers = HeaderMap::new(); + assert!(upload_offset(&headers).is_err(), "missing header"); + + headers.insert(HEADER_UPLOAD_OFFSET, "*".parse().unwrap()); + assert_eq!(upload_offset(&headers).unwrap(), UploadOffset::Unknown); + + headers.insert(HEADER_UPLOAD_OFFSET, "262144".parse().unwrap()); + assert_eq!(upload_offset(&headers).unwrap(), UploadOffset::At(262_144)); + + headers.insert(HEADER_UPLOAD_OFFSET, "nope".parse().unwrap()); + assert!(upload_offset(&headers).is_err()); + } + + /// 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(); + + (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()).unwrap(); + + let (status, offset, body) = parts_of(response).await; + assert_eq!(status, StatusCode::NO_CONTENT); + assert_eq!(offset.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::Committed), "my-key".into()).unwrap(); + + let (status, offset, body) = parts_of(response).await; + assert_eq!(status, StatusCode::CREATED); + assert_eq!(offset, None, "a commit reports no offset"); + assert_eq!(body, r#"{"key":"my-key"}"#); + } + + #[tokio::test] + async fn offset_mismatch_answers_conflict_with_the_authoritative_offset() { + let mismatch = ServiceError::UploadOffsetMismatch { offset: 786_432 }; + let response = + progress_response(Err(ApiError::Service(mismatch)), "my-key".into()).unwrap(); + + 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!(body.contains("786432"), "{body:?}"); + } + + #[tokio::test] + async fn other_errors_propagate_unchanged() { + let gone = ApiError::Service(ServiceError::UploadSessionGone); + let error = progress_response(Err(gone), "my-key".into()).unwrap_err(); + assert_eq!(error.status(), StatusCode::GONE); + + let oversized = ApiError::Service(ServiceError::ChunkExceedsUploadLength { + offset: 8, + content_length: 4, + upload_length: 10, + }); + let error = progress_response(Err(oversized), "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 new file mode 100644 index 00000000..8a7a6592 --- /dev/null +++ b/objectstore-server/tests/resumable.rs @@ -0,0 +1,322 @@ +//! Integration tests for the resumable upload endpoints. +//! +//! No backend implements resumable uploads yet. Until a supporting test backend exists, these +//! tests cover request validation and ensure regular object requests remain unaffected. + +use std::io::{Read, Write}; +use std::net::TcpStream; + +use anyhow::Result; +use objectstore_server::config::{AuthZ, Config}; +use objectstore_test::server::TestServer; +use objectstore_types::resumable::{HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET}; +use reqwest::StatusCode; + +/// Unpadded base64url for the opaque backend token `some-token`. +const SESSION: &str = "c29tZS10b2tlbg"; + +async fn test_server() -> TestServer { + TestServer::with_config(Config { + auth: AuthZ { + enforce: false, + ..Default::default() + }, + ..Default::default() + }) + .await +} + +/// Sends a raw HTTP/1.1 `PUT`, preserving the caller's exact body framing headers. +async fn raw_put(server: &TestServer, path: &str, headers: &str, body: &str) -> Result { + let url = reqwest::Url::parse(&server.url(path))?; + let host = url + .host_str() + .expect("test server URL has a host") + .to_owned(); + let port = url + .port_or_known_default() + .expect("test server URL has a port"); + let target = match url.query() { + Some(query) => format!("{}?{query}", url.path()), + None => url.path().to_owned(), + }; + let headers = headers.to_owned(); + let body = body.to_owned(); + + tokio::task::spawn_blocking(move || -> Result { + let mut stream = TcpStream::connect((host.as_str(), port))?; + write!( + stream, + "PUT {target} HTTP/1.1\r\nHost: {host}:{port}\r\n{headers}Connection: close\r\n\r\n{body}" + )?; + + let mut response = String::new(); + stream.read_to_string(&mut response)?; + Ok(response) + }) + .await? +} + +// --- Session creation --- + +#[tokio::test] +async fn create_session_requires_upload_length() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url("/v1/objects/test/org=1/my-key?upload_type=resumable")) + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +#[tokio::test] +async fn create_session_rejects_malformed_upload_length() -> Result<()> { + let server = test_server().await; + let client = reqwest::Client::new(); + + for invalid in ["", "-1", "1.5", "lots"] { + let response = client + .put(server.url("/v1/objects/test/org=1/my-key?upload_type=resumable")) + .header(HEADER_UPLOAD_LENGTH, invalid) + .send() + .await?; + + assert_eq!( + response.status(), + StatusCode::BAD_REQUEST, + "accepted {HEADER_UPLOAD_LENGTH}: {invalid:?}" + ); + } + + Ok(()) +} + +#[tokio::test] +async fn create_session_rejects_a_payload() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url("/v1/objects/test/org=1/my-key?upload_type=resumable")) + .header(HEADER_UPLOAD_LENGTH, "7") + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +#[tokio::test] +async fn unknown_upload_type_is_rejected() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url("/v1/objects/test/org=1/my-key?upload_type=multipart")) + .header(HEADER_UPLOAD_LENGTH, "1048576") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +// --- Chunks and offset queries --- + +#[tokio::test] +async fn offset_query_rejects_chunked_body_without_content_length() -> Result<()> { + let server = test_server().await; + let response = raw_put( + &server, + &format!("/v1/objects/test/org=1/my-key?session={SESSION}"), + &format!("{HEADER_UPLOAD_OFFSET}: *\r\nTransfer-Encoding: chunked\r\n"), + "7\r\npayload\r\n0\r\n\r\n", + ) + .await?; + + assert!( + response.starts_with("HTTP/1.1 400 Bad Request\r\n"), + "unexpected response: {response}" + ); + Ok(()) +} + +#[tokio::test] +async fn chunk_requires_upload_offset() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url(&format!("/v1/objects/test/org=1/my-key?session={SESSION}"))) + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +#[tokio::test] +async fn chunk_rejects_malformed_upload_offset() -> Result<()> { + let server = test_server().await; + let client = reqwest::Client::new(); + + for invalid in ["", "-1", "1.5", "**", "here"] { + let response = client + .put(server.url(&format!("/v1/objects/test/org=1/my-key?session={SESSION}"))) + .header(HEADER_UPLOAD_OFFSET, invalid) + .body("payload") + .send() + .await?; + + assert_eq!( + response.status(), + StatusCode::BAD_REQUEST, + "accepted {HEADER_UPLOAD_OFFSET}: {invalid:?}" + ); + } + + Ok(()) +} + +#[tokio::test] +async fn offset_query_rejects_a_payload() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url(&format!("/v1/objects/test/org=1/my-key?session={SESSION}"))) + .header(HEADER_UPLOAD_OFFSET, "*") + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +#[tokio::test] +async fn session_token_requires_base64url() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url("/v1/objects/test/org=1/my-key?session=%25%25%25")) + .header(HEADER_UPLOAD_OFFSET, "0") + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +// --- Cancellation --- + +#[tokio::test] +async fn delete_rejects_upload_type() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .delete(server.url("/v1/objects/test/org=1/my-key?upload_type=resumable")) + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +// --- Parameter combinations --- + +#[tokio::test] +async fn upload_type_and_session_are_mutually_exclusive() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url(&format!( + "/v1/objects/test/org=1/my-key?upload_type=resumable&session={SESSION}" + ))) + .header(HEADER_UPLOAD_LENGTH, "1048576") + .header(HEADER_UPLOAD_OFFSET, "0") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +#[tokio::test] +async fn session_on_the_collection_route_is_rejected() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .post(server.url(&format!("/v1/objects/test/org=1/?session={SESSION}"))) + .header(HEADER_UPLOAD_OFFSET, "0") + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + Ok(()) +} + +// --- Regular uploads are unaffected --- + +#[tokio::test] +async fn regular_object_operations_still_work() -> Result<()> { + let server = test_server().await; + let client = reqwest::Client::new(); + + let response = client + .put(server.url("/v1/objects/test/org=1/my-key")) + .body("payload") + .send() + .await?; + assert_eq!(response.status(), StatusCode::OK); + + let response = client + .get(server.url("/v1/objects/test/org=1/my-key")) + .send() + .await?; + assert_eq!(response.status(), StatusCode::OK); + assert_eq!(response.text().await?, "payload"); + + let response = client + .delete(server.url("/v1/objects/test/org=1/my-key")) + .send() + .await?; + assert_eq!(response.status(), StatusCode::NO_CONTENT); + + Ok(()) +} + +#[tokio::test] +async fn regular_upload_ignores_resumable_headers() -> Result<()> { + let server = test_server().await; + + // Without a query parameter the request is a regular upload, and the protocol headers + // carry no meaning. They must not accidentally engage the resumable path. + let response = reqwest::Client::new() + .put(server.url("/v1/objects/test/org=1/my-key")) + .header(HEADER_UPLOAD_LENGTH, "7") + .header(HEADER_UPLOAD_OFFSET, "0") + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::OK); + Ok(()) +} + +#[tokio::test] +async fn regular_upload_ignores_unrelated_query_parameters() -> Result<()> { + let server = test_server().await; + + let response = reqwest::Client::new() + .put(server.url("/v1/objects/test/org=1/my-key?unrelated=value")) + .body("payload") + .send() + .await?; + + assert_eq!(response.status(), StatusCode::OK); + Ok(()) +} diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 0b745d46..09246076 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -190,6 +190,30 @@ The default execution limit is [`DEFAULT_CONCURRENCY_LIMIT`](service::DEFAULT_CONCURRENCY_LIMIT). See [`StorageService::with_concurrency`] for configuration. +## Resumable Uploads + +A resumable upload writes one object across several requests: the payload arrives +as a sequence of chunks at increasing byte offsets, so an interrupted transfer +continues where it stopped instead of starting over. That is worth the extra round +trips for objects large enough that re-sending the whole payload is expensive. + +[`StorageService`] exposes four operations, each a method on +[`Backend`](backend::common::Backend), which run in this sequence: + +1. [`create_upload_session`](backend::common::Backend::create_upload_session) + declares the total size and metadata, and returns a session token. +2. [`put_chunk`](backend::common::Backend::put_chunk) writes bytes at an offset and + reports the offset now persisted. +3. After a failure, [`upload_offset`](backend::common::Backend::upload_offset) + reports where the backend stands, so the caller resumes from there. +4. The chunk carrying the last byte commits the object. There is no finalize call — + the backend recognizes that chunk from the declared total size. +5. At any time, an upload can be canceled, which discards what its session holds. + +Not all backends support resumable uploads. A backend that refuses to create a session +returns `NotImplemented`; support can depend on the declared size, the metadata, or whether +resuming is possible in principle. + ## Multipart Uploads When the configured backend supports it, [`StorageService`] exposes multipart diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index 45b9b4a0..69c8c771 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -13,6 +13,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; +use crate::resumable::{SessionToken, UploadProgress}; use crate::stream::{ClientStream, PayloadStream}; /// User agent string used for outgoing requests. @@ -72,6 +73,56 @@ pub trait Backend: fmt::Debug + Send + Sync + 'static { fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> { Err(Error::NotImplemented) } + + /// Creates a resumable upload session for the object at `id`. + /// + /// Object metadata and its total length are declared upfront and cannot be mutated + /// during the upload. + /// + /// Returns [`Error::NotImplemented`] when the backend refuses to create a session. + async fn create_upload_session( + &self, + id: &ObjectId, + metadata: &Metadata, + total_length: u64, + ) -> Result { + let _ = (id, metadata, total_length); + Err(Error::NotImplemented) + } + + /// Writes a chunk of `content_length` bytes at `offset` into an open session. + /// + /// `offset` must equal the offset the backend currently holds. + /// Returns [`Error::UnknownUploadSession`] when `session` does not identify an open session, + /// and [`Error::ChunkExceedsUploadLength`] when the chunk would exceed the total length + /// declared when the session was created. + async fn put_chunk( + &self, + id: &ObjectId, + session: &SessionToken, + offset: u64, + content_length: u64, + stream: ClientStream, + ) -> Result { + let _ = (id, session, offset, content_length, stream); + Err(Error::NotImplemented) + } + + /// Reports how far the session has progressed. + /// + /// Returns [`Error::UnknownUploadSession`] when `session` does not identify an open session. + async fn upload_offset(&self, id: &ObjectId, session: &SessionToken) -> Result { + let _ = (id, session); + Err(Error::NotImplemented) + } + + /// Cancels an upload session, discarding whatever was uploaded. + /// + /// Returns [`Error::UnknownUploadSession`] when `session` does not identify an open session. + async fn cancel_upload(&self, id: &ObjectId, session: &SessionToken) -> Result<()> { + let _ = (id, session); + Err(Error::NotImplemented) + } } /// Trait for backends that support our S3-style multipart upload protocol. diff --git a/objectstore-service/src/backend/counting.rs b/objectstore-service/src/backend/counting.rs index 4e597e0f..5ebcc6b5 100644 --- a/objectstore-service/src/backend/counting.rs +++ b/objectstore-service/src/backend/counting.rs @@ -28,6 +28,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; +use crate::resumable::{SessionToken, UploadProgress}; use crate::stream::ClientStream; /// Increments `cogs.usage` by one operation for the given `usecase`. @@ -99,6 +100,42 @@ impl Backend for CountingBackend { self.inner.as_multipart_upload_backend()?; Ok(self) } + + async fn create_upload_session( + &self, + id: &ObjectId, + metadata: &Metadata, + total_length: u64, + ) -> Result { + count(&id.context.usecase); + self.inner + .create_upload_session(id, metadata, total_length) + .await + } + + async fn put_chunk( + &self, + id: &ObjectId, + session: &SessionToken, + offset: u64, + content_length: u64, + stream: ClientStream, + ) -> Result { + count(&id.context.usecase); + self.inner + .put_chunk(id, session, offset, content_length, stream) + .await + } + + async fn upload_offset(&self, id: &ObjectId, session: &SessionToken) -> Result { + count(&id.context.usecase); + self.inner.upload_offset(id, session).await + } + + async fn cancel_upload(&self, id: &ObjectId, session: &SessionToken) -> Result<()> { + count(&id.context.usecase); + self.inner.cancel_upload(id, session).await + } } #[async_trait::async_trait] diff --git a/objectstore-service/src/backend/testing.rs b/objectstore-service/src/backend/testing.rs index 24034f24..ac2573cc 100644 --- a/objectstore-service/src/backend/testing.rs +++ b/objectstore-service/src/backend/testing.rs @@ -52,6 +52,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; +use crate::resumable::{SessionToken, UploadProgress}; use crate::stream::ClientStream; /// Hooks for [`TestBackend`]. @@ -237,6 +238,56 @@ pub trait Hooks: fmt::Debug + Send + Sync + 'static { ) -> Result { inner.complete_multipart(id, upload_id, parts).await } + + // --- Resumable upload methods --- + + /// Intercepts [`Backend::create_upload_session`]. Default delegates to `inner`. + async fn create_upload_session( + &self, + inner: &InMemoryBackend, + id: &ObjectId, + metadata: &Metadata, + total_length: u64, + ) -> Result { + inner + .create_upload_session(id, metadata, total_length) + .await + } + + /// Intercepts [`Backend::put_chunk`]. Default delegates to `inner`. + async fn put_chunk( + &self, + inner: &InMemoryBackend, + id: &ObjectId, + session: &SessionToken, + offset: u64, + content_length: u64, + stream: ClientStream, + ) -> Result { + inner + .put_chunk(id, session, offset, content_length, stream) + .await + } + + /// Intercepts [`Backend::upload_offset`]. Default delegates to `inner`. + async fn upload_offset( + &self, + inner: &InMemoryBackend, + id: &ObjectId, + session: &SessionToken, + ) -> Result { + inner.upload_offset(id, session).await + } + + /// Intercepts [`Backend::cancel_upload`]. Default delegates to `inner`. + async fn cancel_upload( + &self, + inner: &InMemoryBackend, + id: &ObjectId, + session: &SessionToken, + ) -> Result<()> { + inner.cancel_upload(id, session).await + } } /// Generic test backend that implements both [`Backend`] and [`HighVolumeBackend`]. @@ -311,6 +362,38 @@ impl Backend for TestBackend { async fn join(&self) { self.hooks.join(&self.inner).await } + + async fn create_upload_session( + &self, + id: &ObjectId, + metadata: &Metadata, + total_length: u64, + ) -> Result { + self.hooks + .create_upload_session(&self.inner, id, metadata, total_length) + .await + } + + async fn put_chunk( + &self, + id: &ObjectId, + session: &SessionToken, + offset: u64, + content_length: u64, + stream: ClientStream, + ) -> Result { + self.hooks + .put_chunk(&self.inner, id, session, offset, content_length, stream) + .await + } + + async fn upload_offset(&self, id: &ObjectId, session: &SessionToken) -> Result { + self.hooks.upload_offset(&self.inner, id, session).await + } + + async fn cancel_upload(&self, id: &ObjectId, session: &SessionToken) -> Result<()> { + self.hooks.cancel_upload(&self.inner, id, session).await + } } #[async_trait::async_trait] diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index a6af1342..65ba32be 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -96,6 +96,18 @@ //! already-mutated state and still returns `true` — so callers do not mistakenly //! treat a successful commit as a lost race and clean up data that was actually //! persisted. +//! +//! # Resumable Uploads +//! +//! Not implemented here yet, so [`TieredStorage`] inherits the unsupported defaults from +//! [`Backend`] and every session creation returns [`Error::NotImplemented`]. 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 [`Error::NotImplemented`] 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. use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; diff --git a/objectstore-service/src/error.rs b/objectstore-service/src/error.rs index 9643dd43..3f8dd71d 100644 --- a/objectstore-service/src/error.rs +++ b/objectstore-service/src/error.rs @@ -151,6 +151,35 @@ pub enum Error { /// Invalid upload ID (e.g. path traversal attempt). #[error(transparent)] InvalidUploadId(#[from] objectstore_types::multipart::InvalidUploadId), + + /// A resumable chunk was submitted at an offset that's different from the one held by the + /// backend. + #[error("upload offset mismatch (server holds {offset} bytes)")] + UploadOffsetMismatch { + /// The offset the backend currently holds. + offset: u64, + }, + + /// The resumable upload session expired or was canceled, retaining nothing. + #[error("upload session gone")] + UploadSessionGone, + + /// The backend does not recognize the addressed resumable upload session. + #[error("unknown upload session")] + UnknownUploadSession, + + /// A resumable upload chunk would exceed the length declared for the session. + #[error( + "chunk at offset {offset} with length {content_length} exceeds upload length {upload_length}" + )] + ChunkExceedsUploadLength { + /// The offset at which the chunk would be written. + offset: u64, + /// The declared length of the chunk. + content_length: u64, + /// The total upload length declared when the session was created. + upload_length: u64, + }, } impl Error { @@ -197,6 +226,14 @@ impl Error { Self::Client(_) => Level::DEBUG, Self::Metadata(_) => Level::DEBUG, Self::RangeNotSatisfiable { .. } => Level::DEBUG, + Self::UploadOffsetMismatch { .. } => Level::DEBUG, + Self::UploadSessionGone => Level::DEBUG, + Self::UnknownUploadSession => Level::DEBUG, + Self::ChunkExceedsUploadLength { .. } => 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. + Self::NotImplemented => Level::INFO, // Like rate limits, we treat capacity errors as warnings Self::AtCapacity => Level::WARN, // All other errors are service or backend failures @@ -208,7 +245,6 @@ impl Error { Self::Panic(_) => Level::ERROR, Self::Dropped => Level::ERROR, Self::UnexpectedTombstone => Level::ERROR, - Self::NotImplemented => Level::ERROR, Self::InvalidUploadId(_) => Level::DEBUG, Self::Generic { .. } => Level::ERROR, } diff --git a/objectstore-service/src/lib.rs b/objectstore-service/src/lib.rs index e33a2f7c..d1a2d364 100644 --- a/objectstore-service/src/lib.rs +++ b/objectstore-service/src/lib.rs @@ -8,6 +8,7 @@ pub mod error; mod gcp_auth; pub mod id; pub mod multipart; +pub mod resumable; pub mod service; pub mod stream; pub mod streaming; diff --git a/objectstore-service/src/resumable.rs b/objectstore-service/src/resumable.rs new file mode 100644 index 00000000..3aba2f7d --- /dev/null +++ b/objectstore-service/src/resumable.rs @@ -0,0 +1,36 @@ +//! Shared types for Objectstore's resumable upload protocol. +//! +//! A resumable upload is a regular write whose payload arrives across several requests. +//! A session declares the object's total size and metadata upfront; chunks then arrive at +//! increasing byte offsets, and the backend commits the object itself once the last byte +//! lands. See [`objectstore_types::resumable`] for the wire-level types. +//! +//! Not every backend can support this. Session creation therefore asks the backend that +//! would store the object to open one. A backend that cannot create a session returns +//! [`Error::NotImplemented`](crate::error::Error::NotImplemented), and the client falls back to +//! a regular upload. + +pub use objectstore_types::resumable::{SessionToken, UploadOffset}; + +/// How far a resumable upload has progressed. +/// +/// Returned by both +/// [`Backend::put_chunk`](crate::backend::common::Backend::put_chunk) and +/// [`Backend::upload_offset`](crate::backend::common::Backend::upload_offset), because an +/// offset query commits an object that was assembled but not yet committed and therefore +/// has the same two outcomes as a chunk write. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum UploadProgress { + /// More bytes are expected. The client continues from `offset`. + /// + /// This offset is authoritative and may be lower than the end of the chunk that was + /// just written: backends persist only aligned prefixes and discard the remainder. It must + /// remain below the session's total length; once every byte has landed, the backend commits + /// the object or returns an error instead. + Incomplete { + /// The offset the backend has persisted. + offset: u64, + }, + /// The last byte arrived and the object is committed and readable. + Committed, +} diff --git a/objectstore-service/src/service.rs b/objectstore-service/src/service.rs index 0953fcd8..95fbe2b9 100644 --- a/objectstore-service/src/service.rs +++ b/objectstore-service/src/service.rs @@ -20,6 +20,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; +use crate::resumable::{SessionToken, UploadProgress}; use crate::stream::{ClientStream, PayloadStream}; use crate::streaming::StreamExecutor; @@ -359,6 +360,78 @@ impl StorageService { }) .await } + + // --- Resumable upload operations --- + + /// Opens a resumable upload session for an object of `total_length` bytes. + /// + /// Returns [`Error::NotImplemented`](crate::error::Error::NotImplemented) when the backend + /// refuses to create a session, in which case the caller should fall back to + /// [`Self::insert_object`]. + pub async fn create_upload_session( + &self, + id: ObjectId, + metadata: Metadata, + total_length: u64, + ) -> Result { + metadata.validate()?; + let inner = Arc::clone(&self.inner); + self.spawn("create_upload_session", async move { + inner + .create_upload_session(&id, &metadata, total_length) + .await + }) + .await + } + + /// Writes a chunk of `content_length` bytes at `offset` into an open session. + /// + /// Commits the object once the chunk carrying the last byte is persisted. + /// + /// # Run-to-completion + /// + /// Once called, the operation runs to completion even if the returned future is dropped. + /// This matters most for the final chunk, which commits the object. + pub async fn put_chunk( + &self, + id: ObjectId, + session: SessionToken, + offset: u64, + content_length: u64, + body: ClientStream, + ) -> Result { + let inner = Arc::clone(&self.inner); + self.spawn("put_chunk", async move { + inner + .put_chunk(&id, &session, offset, content_length, body) + .await + }) + .await + } + + /// Reports how far a session has progressed, committing the object if it is assembled. + /// + /// This can mutate state and therefore requires write permission at the API layer. + pub async fn upload_offset( + &self, + id: ObjectId, + session: SessionToken, + ) -> Result { + let inner = Arc::clone(&self.inner); + self.spawn("upload_offset", async move { + inner.upload_offset(&id, &session).await + }) + .await + } + + /// Cancels an upload session, discarding whatever was uploaded. + pub async fn cancel_upload(&self, id: ObjectId, session: SessionToken) -> Result<()> { + let inner = Arc::clone(&self.inner); + self.spawn("cancel_upload", async move { + inner.cancel_upload(&id, &session).await + }) + .await + } } #[cfg(test)] @@ -368,7 +441,7 @@ mod tests { use bytes::BytesMut; use futures_util::TryStreamExt; - use objectstore_types::metadata::Metadata; + use objectstore_types::metadata::{ExpirationPolicy, Metadata}; use objectstore_types::range::ByteRange; use objectstore_types::scope::{Scope, Scopes}; @@ -781,4 +854,22 @@ mod tests { "permit was not released after panic" ); } + + // --- Resumable uploads --- + + #[tokio::test] + async fn resumable_create_validates_metadata() { + let service = make_service(); + let id = ObjectId::new(make_context(), "resumable".into()); + + // A timeout policy with no resolved `time_expires` is rejected before the backend + // is consulted, exactly as it is for a regular insert. + let metadata = Metadata { + expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_secs(60)), + ..Default::default() + }; + + let result = service.create_upload_session(id, metadata, 1024).await; + assert!(matches!(result, Err(Error::Metadata(_))), "{result:?}"); + } } diff --git a/objectstore-types/src/lib.rs b/objectstore-types/src/lib.rs index 910f293e..79ce8870 100644 --- a/objectstore-types/src/lib.rs +++ b/objectstore-types/src/lib.rs @@ -11,4 +11,5 @@ pub mod metadata; pub mod multipart; pub mod presign; pub mod range; +pub mod resumable; pub mod scope; diff --git a/objectstore-types/src/resumable.rs b/objectstore-types/src/resumable.rs new file mode 100644 index 00000000..48697aaa --- /dev/null +++ b/objectstore-types/src/resumable.rs @@ -0,0 +1,140 @@ +//! Types for the resumable upload protocol. + +use std::fmt; +use std::str::FromStr; + +use serde::{Deserialize, Serialize}; + +/// Request header declaring the total size of the object, in bytes. +/// +/// Required when creating a session. +pub const HEADER_UPLOAD_LENGTH: &str = "upload-length"; + +/// Header carrying the byte offset of a chunk, or the offset the server holds. +/// +/// On a request this is the offset of the chunk's first byte, or `*` to query the +/// server's authoritative offset. On a response it is the offset the server has +/// persisted. See [`UploadOffset`]. +pub const HEADER_UPLOAD_OFFSET: &str = "upload-offset"; + +/// The wildcard [`HEADER_UPLOAD_OFFSET`] value that queries the server's offset. +const OFFSET_WILDCARD: &str = "*"; + +/// Identifier for an in-progress resumable upload session. +/// +/// The token is an opaque identifier whose contents are defined and interpreted by the storage +/// backend. +pub type SessionToken = String; + +/// The value of the [`HEADER_UPLOAD_OFFSET`] request header. +/// +/// In a request, a concrete offset submits a chunk starting at that byte, +/// while [`UploadOffset::Unknown`] asks the server which offset it holds. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum UploadOffset { + /// Denotes a chunk whose first byte sits at this offset. + At(u64), + /// Used to query the server for its authoritative offset. + Unknown, +} + +/// Error returned when an [`UploadOffset`] header value cannot be parsed. +#[derive(Debug, thiserror::Error)] +#[error("invalid {HEADER_UPLOAD_OFFSET} value: {0}")] +pub struct InvalidUploadOffset(String); + +impl FromStr for UploadOffset { + type Err = InvalidUploadOffset; + + fn from_str(s: &str) -> Result { + if s == OFFSET_WILDCARD { + return Ok(Self::Unknown); + } + + // Rejects the `+` sign and leading whitespace that `u64::from_str` would + // otherwise be lenient about, keeping the header canonical. + if !s.bytes().all(|b| b.is_ascii_digit()) { + return Err(InvalidUploadOffset(s.to_owned())); + } + + let offset = s.parse().map_err(|_| InvalidUploadOffset(s.to_owned()))?; + Ok(Self::At(offset)) + } +} + +impl fmt::Display for UploadOffset { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::At(offset) => offset.fmt(f), + Self::Unknown => f.write_str(OFFSET_WILDCARD), + } + } +} + +/// Response from creating a resumable upload session. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CreateSessionResponse { + /// The object key (server-generated or client-provided). + pub key: String, + /// The opaque session token for subsequent requests. + pub session: SessionToken, +} + +/// Response from the request that commits the object. +/// +/// This is either the chunk carrying the last byte, or an offset query against a +/// session whose object was already assembled. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CommitResponse { + /// The object key. + pub key: String, +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn create_session_response_serializes_token_verbatim() -> Result<(), serde_json::Error> { + let response = CreateSessionResponse { + key: "key".into(), + session: "../opaque +? ü".into(), + }; + + assert_eq!( + serde_json::to_string(&response)?, + r#"{"key":"key","session":"../opaque +? ü"}"# + ); + Ok(()) + } + + #[test] + fn upload_offset_parses_wildcard_and_offsets() -> Result<(), InvalidUploadOffset> { + assert_eq!("*".parse::()?, UploadOffset::Unknown); + assert_eq!("0".parse::()?, UploadOffset::At(0)); + assert_eq!("262144".parse::()?, UploadOffset::At(262144)); + Ok(()) + } + + #[test] + fn upload_offset_rejects_malformed_values() { + for invalid in ["", "-1", "+1", " 1", "1 ", "1.5", "0x10", "**", "abc"] { + assert!( + invalid.parse::().is_err(), + "expected {invalid:?} to be rejected" + ); + } + } + + #[test] + fn upload_offset_round_trips_through_display() -> Result<(), InvalidUploadOffset> { + for offset in [ + UploadOffset::Unknown, + UploadOffset::At(0), + UploadOffset::At(7), + ] { + assert_eq!(offset.to_string().parse::()?, offset); + } + Ok(()) + } +}