From 1643654f3d55cd8aeae80b90f0eabec198fd7a9a Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Tue, 8 Sep 2026 15:40:19 +0200 Subject: [PATCH] ref(service): Standardize backend session token codecs Give concrete backends typed session-token codecs while keeping resumable methods object-safe and preserving the opaque wire format. --- objectstore-service/src/backend/bigtable.rs | 2 + objectstore-service/src/backend/common.rs | 25 +++++++- objectstore-service/src/backend/counting.rs | 2 + objectstore-service/src/backend/gcs.rs | 57 ++++++++++--------- objectstore-service/src/backend/in_memory.rs | 2 + objectstore-service/src/backend/local_fs.rs | 2 + .../src/backend/s3_compatible.rs | 2 + objectstore-service/src/backend/testing.rs | 2 + objectstore-service/src/backend/tiered.rs | 2 + objectstore-service/src/resumable.rs | 16 +++++- objectstore-service/src/service.rs | 7 ++- 11 files changed, 89 insertions(+), 30 deletions(-) diff --git a/objectstore-service/src/backend/bigtable.rs b/objectstore-service/src/backend/bigtable.rs index d0b0b795..5246ec4a 100644 --- a/objectstore-service/src/backend/bigtable.rs +++ b/objectstore-service/src/backend/bigtable.rs @@ -987,6 +987,8 @@ impl BigTableBackend { #[async_trait::async_trait] impl Backend for BigTableBackend { + type SessionToken = (); + fn name(&self) -> &'static str { "bigtable" } diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index 2f8c17d0..3b4b7bd7 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -36,6 +36,29 @@ pub type DeleteResponse = (); /// Trait implemented by all storage backends. #[async_trait::async_trait] pub trait Backend: fmt::Debug + Send + Sync + 'static { + /// Backend-specific resumable upload session state. + type SessionToken + where + Self: Sized; + + /// Encodes backend-specific session state into an opaque token. + fn encode_session_token(token: Self::SessionToken) -> Result + where + Self: Sized, + { + let _ = token; + Err(ErrorKind::Unsupported.into()) + } + + /// Decodes an opaque token into backend-specific session state. + fn decode_session_token(token: &BackendToken) -> Result + where + Self: Sized, + { + let _ = token; + Err(ErrorKind::Unsupported.into()) + } + /// The backend name, used for diagnostics. fn name(&self) -> &'static str; @@ -92,7 +115,7 @@ pub trait Backend: fmt::Debug + Send + Sync + 'static { /// Object metadata and its total length are declared upfront and cannot be mutated /// during the upload. /// - /// The returned string is opaque backend-defined state. [`StorageService`](crate::StorageService) + /// The returned token contains opaque backend-defined state. [`StorageService`](crate::StorageService) /// protects it before exposing the session token outside the service layer. /// /// Returns `Ok(None)` when this backend cannot store the described object resumably. Declining diff --git a/objectstore-service/src/backend/counting.rs b/objectstore-service/src/backend/counting.rs index 993189da..87fac591 100644 --- a/objectstore-service/src/backend/counting.rs +++ b/objectstore-service/src/backend/counting.rs @@ -66,6 +66,8 @@ impl CountingBackend { #[async_trait::async_trait] impl Backend for CountingBackend { + type SessionToken = (); + fn name(&self) -> &'static str { self.inner.name() } diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 7646143c..6ef4e61c 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -486,7 +486,7 @@ const CLIENT_CLOSED_REQUEST_STATUS: u16 = 499; /// Represents a resumable upload session in GCS. #[derive(Debug)] -struct ResumableUpload { +pub struct GcsSessionToken { // URI to use for requests that act on this session, returned by GCS in the `Location` header // on session creation. session_uri: Url, @@ -494,29 +494,13 @@ struct ResumableUpload { total_length: NonZeroU64, } -impl ResumableUpload { +impl GcsSessionToken { fn new(session_uri: Url, total_length: NonZeroU64) -> Self { Self { session_uri, total_length, } } - - fn into_token(self) -> BackendToken { - format!("{}.{}", self.total_length, self.session_uri) - } - - fn from_token(token: &BackendToken) -> Result { - let (total_length, session_uri) = token - .split_once('.') - .ok_or(ErrorKind::UnknownUploadSession)?; - let total_length = total_length - .parse::() - .map_err(|_| ErrorKind::UnknownUploadSession)?; - let session_uri = Url::parse(session_uri).map_err(|_| ErrorKind::UnknownUploadSession)?; - let session = Self::new(session_uri, total_length); - Ok(session) - } } enum GcsUploadProgress { @@ -830,7 +814,7 @@ fn range_header_to_offset(value: &str, total_length: NonZeroU64) -> Result /// Returns the progress GCS reported, plus the completed object when this is the response that /// finished the upload. async fn range_response_to_upload_progress( - session: &ResumableUpload, + session: &GcsSessionToken, response: reqwest::Response, ) -> Result { let status = response.status(); @@ -893,6 +877,27 @@ async fn range_response_to_upload_progress( #[async_trait::async_trait] impl Backend for GcsBackend { + type SessionToken = GcsSessionToken; + + fn encode_session_token(token: Self::SessionToken) -> Result { + Ok(BackendToken::new(format!( + "{}.{}", + token.total_length, token.session_uri + ))) + } + + fn decode_session_token(token: &BackendToken) -> Result { + let (total_length, session_uri) = token + .as_str() + .split_once('.') + .ok_or(ErrorKind::UnknownUploadSession)?; + let total_length = total_length + .parse::() + .map_err(|_| ErrorKind::UnknownUploadSession)?; + let session_uri = Url::parse(session_uri).map_err(|_| ErrorKind::UnknownUploadSession)?; + Ok(GcsSessionToken::new(session_uri, total_length)) + } + fn name(&self) -> &'static str { "gcs" } @@ -1186,8 +1191,8 @@ impl Backend for GcsBackend { "invalid Location URL in GCS resumable upload creation response", ) })?; - let session = ResumableUpload::new(session_uri, total_length); - Ok(Some(session.into_token())) + let session = GcsSessionToken::new(session_uri, total_length); + Ok(Some(Self::encode_session_token(session)?)) } #[tracing::instrument(level = "debug", fields(?id, offset, content_length), skip_all)] @@ -1199,8 +1204,8 @@ impl Backend for GcsBackend { content_length: u64, stream: ClientStream, ) -> Result { + let session = Self::decode_session_token(token)?; objectstore_log::debug!("Uploading resumable chunk to GCS backend"); - let session = ResumableUpload::from_token(token)?; let end = offset .checked_add(content_length) @@ -1237,8 +1242,8 @@ impl Backend for GcsBackend { #[tracing::instrument(level = "debug", fields(?id), skip_all)] async fn upload_offset(&self, id: &ObjectId, token: &BackendToken) -> Result { + let session = Self::decode_session_token(token)?; objectstore_log::debug!("Querying resumable upload offset on GCS backend"); - let session = ResumableUpload::from_token(token)?; self.with_retry("query_resumable_upload", || async { let response = self @@ -1271,8 +1276,8 @@ impl Backend for GcsBackend { #[tracing::instrument(level = "debug", fields(?id), skip_all)] async fn cancel_upload(&self, id: &ObjectId, token: &BackendToken) -> Result<()> { + let session = Self::decode_session_token(token)?; objectstore_log::debug!("Cancelling resumable upload on GCS backend"); - let session = ResumableUpload::from_token(token)?; let session_uri = session.session_uri; self.with_retry("cancel_resumable_upload", || { let session_uri = session_uri.clone(); @@ -1900,9 +1905,9 @@ mod tests { #[test] fn resumable_token_rejects_zero_length() { - let token = "0.http://localhost/upload".to_owned(); + let token = BackendToken::new("0.http://localhost/upload".to_owned()); assert!(matches!( - ResumableUpload::from_token(&token), + GcsBackend::decode_session_token(&token), Err(error) if error.kind() == ErrorKind::UnknownUploadSession )); } diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index fe2da8c5..96ed39ea 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -111,6 +111,8 @@ impl InMemoryBackend { #[async_trait::async_trait] impl super::common::Backend for InMemoryBackend { + type SessionToken = (); + fn name(&self) -> &'static str { self.name } diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index d71168c9..da94c491 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -96,6 +96,8 @@ impl LocalFsBackend { #[async_trait::async_trait] impl Backend for LocalFsBackend { + type SessionToken = (); + fn name(&self) -> &'static str { "local-fs" } diff --git a/objectstore-service/src/backend/s3_compatible.rs b/objectstore-service/src/backend/s3_compatible.rs index fd27dceb..c0d81397 100644 --- a/objectstore-service/src/backend/s3_compatible.rs +++ b/objectstore-service/src/backend/s3_compatible.rs @@ -320,6 +320,8 @@ impl S3CompatibleBackend { #[async_trait::async_trait] impl Backend for S3CompatibleBackend { + type SessionToken = (); + fn name(&self) -> &'static str { "s3-compatible" } diff --git a/objectstore-service/src/backend/testing.rs b/objectstore-service/src/backend/testing.rs index 0bc997d1..e80c1210 100644 --- a/objectstore-service/src/backend/testing.rs +++ b/objectstore-service/src/backend/testing.rs @@ -353,6 +353,8 @@ impl TestBackend { #[async_trait::async_trait] impl Backend for TestBackend { + type SessionToken = (); + fn name(&self) -> &'static str { self.hooks.name() } diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 5e80d732..57f4b9a0 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -382,6 +382,8 @@ impl TieredStorage { #[async_trait::async_trait] impl Backend for TieredStorage { + type SessionToken = (); + fn name(&self) -> &'static str { "tiered" } diff --git a/objectstore-service/src/resumable.rs b/objectstore-service/src/resumable.rs index 033ebf82..6521fcee 100644 --- a/objectstore-service/src/resumable.rs +++ b/objectstore-service/src/resumable.rs @@ -19,7 +19,21 @@ pub use objectstore_types::resumable::{ }; /// Opaque session state encoded and decoded by a storage backend. -pub type BackendToken = String; +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +#[serde(transparent)] +pub struct BackendToken(String); + +impl BackendToken { + /// Creates an opaque backend token from its encoded representation. + pub fn new(token: String) -> Self { + Self(token) + } + + /// Returns the encoded token representation. + pub fn as_str(&self) -> &str { + &self.0 + } +} /// Structured token encrypted at the service boundary. #[derive(Deserialize, Serialize)] diff --git a/objectstore-service/src/service.rs b/objectstore-service/src/service.rs index 366ec393..c3ab21c2 100644 --- a/objectstore-service/src/service.rs +++ b/objectstore-service/src/service.rs @@ -572,7 +572,7 @@ mod tests { _metadata: &Metadata, _total_length: NonZeroU64, ) -> Result> { - Ok(Some("backend token".to_owned())) + Ok(Some(BackendToken::new("backend token".to_owned()))) } async fn upload_offset( @@ -581,7 +581,10 @@ mod tests { _id: &ObjectId, token: &BackendToken, ) -> Result { - self.seen_tokens.lock().unwrap().push(token.to_owned()); + self.seen_tokens + .lock() + .unwrap() + .push(token.as_str().to_owned()); Ok(UploadProgress::Incomplete { offset: 0 }) } }