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}