Skip to main content

objectstore_types/
resumable.rs

1//! Types shared by resumable upload clients and servers.
2//!
3//! A resumable upload writes one object across multiple requests. The client first creates a
4//! session, declaring the object's complete size with [`HEADER_UPLOAD_LENGTH`]. The server returns
5//! a [`CreateSessionResponse`] containing an opaque [`SessionToken`] that identifies the upload.
6//! The response also reports the upload granularity. A reconstructed client does not know it.
7//!
8//! The client then sends chunks with [`HEADER_UPLOAD_OFFSET`] set to the byte position at which
9//! each chunk starts. If an upload is interrupted, the client can send the wildcard offset
10//! [`UploadOffset::Unknown`] to query the server's authoritative position before resuming. The
11//! request that completes the upload returns a [`CompleteUploadResponse`].
12//!
13//! Session tokens contain the canonical object path and backend state protected by the storage
14//! service, and clients must treat their contents as opaque. The token bytes are encoded as
15//! unpadded base64url when the token is placed in a request's `session` query parameter.
16
17use std::str::FromStr;
18use std::{borrow::Cow, fmt};
19
20use base64::Engine as _;
21use base64::engine::general_purpose::URL_SAFE_NO_PAD;
22use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
23
24/// Request header declaring the total size of the object, in bytes.
25///
26/// Required when creating a session.
27pub const HEADER_UPLOAD_LENGTH: &str = "upload-length";
28
29/// Header carrying the byte offset of a chunk, or the offset the server holds.
30///
31/// On a request this is the offset of the chunk's first byte, or `*` to query the
32/// server's authoritative offset. On a response it is the offset the server has
33/// persisted. See [`UploadOffset`].
34pub const HEADER_UPLOAD_OFFSET: &str = "upload-offset";
35
36/// The wildcard [`HEADER_UPLOAD_OFFSET`] value that queries the server's offset.
37const OFFSET_WILDCARD: &str = "*";
38
39/// Identifier for an in-progress resumable upload session.
40///
41/// Internally, this is an opaque byte string interpreted by the storage service. At the HTTP API
42/// boundary it serializes as canonical unpadded base64url, so the serialized value can be placed
43/// directly in a subsequent request URL.
44#[derive(Clone, PartialEq, Eq)]
45pub struct SessionToken(Vec<u8>);
46
47impl SessionToken {
48    /// Wraps opaque session-token bytes.
49    pub fn new(bytes: impl Into<Vec<u8>>) -> Self {
50        Self(bytes.into())
51    }
52
53    /// Returns the opaque token bytes.
54    pub fn as_bytes(&self) -> &[u8] {
55        &self.0
56    }
57
58    /// Consumes the token and returns its opaque bytes.
59    pub fn into_bytes(self) -> Vec<u8> {
60        self.0
61    }
62
63    /// Parses the canonical unpadded-base64url representation used by the HTTP API.
64    pub fn from_base64url(encoded: &str) -> Result<Self, InvalidSessionToken> {
65        let bytes = URL_SAFE_NO_PAD
66            .decode(encoded)
67            .map_err(|_| InvalidSessionToken)?;
68        if URL_SAFE_NO_PAD.encode(&bytes) != encoded {
69            return Err(InvalidSessionToken);
70        }
71        Ok(Self(bytes))
72    }
73
74    /// Encodes this token for the HTTP API.
75    pub fn to_base64url(&self) -> String {
76        URL_SAFE_NO_PAD.encode(&self.0)
77    }
78}
79
80impl fmt::Debug for SessionToken {
81    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
82        f.write_str("SessionToken")
83    }
84}
85
86impl Serialize for SessionToken {
87    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
88    where
89        S: Serializer,
90    {
91        serializer.serialize_str(&self.to_base64url())
92    }
93}
94
95impl<'de> Deserialize<'de> for SessionToken {
96    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
97    where
98        D: Deserializer<'de>,
99    {
100        let encoded = Cow::<'static, String>::deserialize(deserializer)?;
101        Self::from_base64url(&encoded).map_err(de::Error::custom)
102    }
103}
104
105/// Error returned for a non-canonical or malformed external session token.
106#[derive(Debug, thiserror::Error)]
107#[error("session token must use unpadded base64url encoding")]
108pub struct InvalidSessionToken;
109
110/// The value of the [`HEADER_UPLOAD_OFFSET`] request header.
111///
112/// In a request, a concrete offset submits a chunk starting at that byte,
113/// while [`UploadOffset::Unknown`] asks the server which offset it holds.
114#[derive(Clone, Copy, Debug, PartialEq, Eq)]
115pub enum UploadOffset {
116    /// Denotes a chunk whose first byte sits at this offset.
117    At(u64),
118    /// Used to query the server for its authoritative offset.
119    Unknown,
120}
121
122/// Error returned when an [`UploadOffset`] header value cannot be parsed.
123#[derive(Debug, thiserror::Error)]
124#[error("invalid {HEADER_UPLOAD_OFFSET} value: {0}")]
125pub struct InvalidUploadOffset(String);
126
127impl FromStr for UploadOffset {
128    type Err = InvalidUploadOffset;
129
130    fn from_str(s: &str) -> Result<Self, Self::Err> {
131        if s == OFFSET_WILDCARD {
132            return Ok(Self::Unknown);
133        }
134
135        // Rejects the `+` sign and leading whitespace that `u64::from_str` would
136        // otherwise be lenient about, keeping the header canonical.
137        if !s.bytes().all(|b| b.is_ascii_digit()) {
138            return Err(InvalidUploadOffset(s.to_owned()));
139        }
140
141        let offset = s.parse().map_err(|_| InvalidUploadOffset(s.to_owned()))?;
142        Ok(Self::At(offset))
143    }
144}
145
146impl fmt::Display for UploadOffset {
147    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
148        match self {
149            Self::At(offset) => offset.fmt(f),
150            Self::Unknown => f.write_str(OFFSET_WILDCARD),
151        }
152    }
153}
154
155/// How far a resumable upload has progressed.
156///
157/// Both a chunk write and an offset query can observe that an upload is complete, so both
158/// operations have the same two outcomes. Completion is relative to the backend handling the
159/// operation: it means the session is terminal and the object is available through that backend's
160/// normal read methods. A backend that composes another backend must finish its own publication
161/// work before returning [`UploadProgress::Complete`].
162#[derive(Clone, Copy, Debug, PartialEq, Eq)]
163pub enum UploadProgress {
164    /// More bytes are expected. The client continues from `offset`.
165    ///
166    /// This offset is authoritative and may be lower than the end of the chunk that was just
167    /// written: backends can persist only a prefix and discard the remainder. It must remain below
168    /// the session's total length; once every byte has landed, the backend completes the upload or
169    /// returns an error instead.
170    Incomplete {
171        /// The offset the backend has persisted.
172        offset: u64,
173    },
174    /// The session is terminal and the object is available through the backend's normal reads.
175    ///
176    /// This is an observable status rather than a one-time event. A later offset query can return
177    /// `Complete` again, for example when the response to the final chunk was lost.
178    Complete,
179}
180
181/// Response from creating a resumable upload session.
182#[derive(Debug, Clone, Serialize, Deserialize)]
183pub struct CreateSessionResponse {
184    /// The object key (server-generated or client-provided).
185    pub key: String,
186    /// The opaque session token that identifies the session.
187    pub session: SessionToken,
188    /// This upload's granularity in bytes, or zero when it imposes none.
189    #[serde(default)]
190    pub granularity: u64,
191}
192
193/// Response from the request that completes the upload.
194///
195/// This is either the chunk carrying the last byte, or an offset query against a
196/// session whose final chunk completed but whose response was not observed.
197#[derive(Debug, Clone, Serialize, Deserialize)]
198pub struct CompleteUploadResponse {
199    /// The object key.
200    pub key: String,
201}
202
203#[cfg(test)]
204mod tests {
205    use super::*;
206
207    #[test]
208    fn create_session_response_encodes_token_once() -> Result<(), serde_json::Error> {
209        let response = CreateSessionResponse {
210            key: "key".into(),
211            session: SessionToken::new(b"../opaque +? \xc3\xbc"),
212            granularity: 262_144,
213        };
214
215        assert_eq!(
216            serde_json::to_string(&response)?,
217            r#"{"key":"key","session":"Li4vb3BhcXVlICs_IMO8","granularity":262144}"#
218        );
219        Ok(())
220    }
221
222    #[test]
223    fn session_token_round_trips_arbitrary_bytes() -> Result<(), serde_json::Error> {
224        let token = SessionToken::new([0, 1, 2, 0xfe, 0xff]);
225        let json = serde_json::to_string(&token)?;
226        assert_eq!(json, r#""AAEC_v8""#);
227        assert_eq!(serde_json::from_str::<SessionToken>(&json)?, token);
228        Ok(())
229    }
230
231    #[test]
232    fn session_token_rejects_noncanonical_encodings() {
233        for invalid in ["%%%", "dG9rM24="] {
234            assert!(
235                SessionToken::from_base64url(invalid).is_err(),
236                "accepted {invalid:?}"
237            );
238        }
239    }
240
241    #[test]
242    fn upload_offset_parses_wildcard_and_offsets() -> Result<(), InvalidUploadOffset> {
243        assert_eq!("*".parse::<UploadOffset>()?, UploadOffset::Unknown);
244        assert_eq!("0".parse::<UploadOffset>()?, UploadOffset::At(0));
245        assert_eq!("262144".parse::<UploadOffset>()?, UploadOffset::At(262144));
246        Ok(())
247    }
248
249    #[test]
250    fn upload_offset_rejects_malformed_values() {
251        for invalid in ["", "-1", "+1", " 1", "1 ", "1.5", "0x10", "**", "abc"] {
252            assert!(
253                invalid.parse::<UploadOffset>().is_err(),
254                "expected {invalid:?} to be rejected"
255            );
256        }
257    }
258
259    #[test]
260    fn upload_offset_round_trips_through_display() -> Result<(), InvalidUploadOffset> {
261        for offset in [
262            UploadOffset::Unknown,
263            UploadOffset::At(0),
264            UploadOffset::At(7),
265        ] {
266            assert_eq!(offset.to_string().parse::<UploadOffset>()?, offset);
267        }
268        Ok(())
269    }
270}