objectstore_client/
resumable.rs1use std::borrow::Cow;
11use std::collections::BTreeMap;
12use std::fmt;
13
14use bytes::Bytes;
15use objectstore_types::metadata::Metadata;
16use objectstore_types::resumable::{
17 CompleteUploadResponse, CreateSessionResponse, HEADER_UPLOAD_LENGTH, HEADER_UPLOAD_OFFSET,
18 UploadOffset,
19};
20use reqwest::{Method, Response, StatusCode};
21use serde::Serialize;
22
23pub use objectstore_types::resumable::{SessionToken, UploadProgress};
24
25use crate::response::ResponseExt as _;
26use crate::{Compression, Error, ExpirationPolicy, ObjectKey, Session};
27
28#[derive(Serialize)]
29#[serde(rename_all = "snake_case")]
30enum UploadType {
31 Resumable,
32}
33
34#[derive(Serialize)]
35struct UploadTypeQuery {
36 upload_type: UploadType,
37}
38
39#[derive(Serialize)]
40struct SessionQuery<'a> {
41 session: &'a SessionToken,
42}
43
44#[derive(Clone, Debug)]
49pub struct ResumableUpload {
50 session: Session,
51 key: ObjectKey,
52 token: SessionToken,
53}
54
55impl Session {
56 pub fn create_upload(&self, object_length: u64) -> CreateResumableUploadBuilder {
66 let metadata = Metadata {
67 expiration_policy: self.scope.usecase().expiration_policy(),
68 ..Default::default()
69 };
70
71 CreateResumableUploadBuilder {
72 session: self.clone(),
73 total_length: object_length,
74 key: None,
75 metadata,
76 }
77 }
78
79 pub fn resume_upload(&self, key: impl Into<ObjectKey>, token: SessionToken) -> ResumableUpload {
84 ResumableUpload {
85 session: self.clone(),
86 key: key.into(),
87 token,
88 }
89 }
90}
91
92impl ResumableUpload {
93 pub fn key(&self) -> &str {
95 &self.key
96 }
97
98 pub fn token(&self) -> &SessionToken {
100 &self.token
101 }
102
103 pub fn progress(&self) -> UploadProgressBuilder {
105 UploadProgressBuilder {
106 upload: self.clone(),
107 }
108 }
109
110 pub fn put(&self, offset: u64, chunk: impl Into<Bytes>) -> PutChunkBuilder {
119 PutChunkBuilder {
120 upload: self.clone(),
121 offset,
122 chunk: chunk.into(),
123 }
124 }
125
126 pub fn cancel(&self) -> CancelUploadBuilder {
128 CancelUploadBuilder {
129 upload: self.clone(),
130 }
131 }
132
133 fn request(&self, method: Method) -> crate::Result<reqwest::RequestBuilder> {
134 Ok(self
135 .session
136 .request(method, &self.key)?
137 .query(&SessionQuery {
138 session: &self.token,
139 }))
140 }
141}
142
143#[derive(Debug)]
145pub struct CreateResumableUploadBuilder {
146 session: Session,
147 total_length: u64,
148 key: Option<ObjectKey>,
149 metadata: Metadata,
150}
151
152impl CreateResumableUploadBuilder {
153 pub fn key(mut self, key: impl Into<ObjectKey>) -> Self {
155 self.key = Some(key.into()).filter(|key| !key.is_empty());
156 self
157 }
158
159 pub fn content_type(mut self, content_type: impl Into<Cow<'static, str>>) -> Self {
161 self.metadata.content_type = content_type.into();
162 self
163 }
164
165 pub fn expiration_policy(mut self, expiration_policy: ExpirationPolicy) -> Self {
167 self.metadata.expiration_policy = expiration_policy;
168 self
169 }
170
171 pub fn compression(mut self, compression: impl Into<Option<Compression>>) -> Self {
179 self.metadata.compression = compression.into();
180 self
181 }
182
183 pub fn origin(mut self, origin: impl Into<String>) -> Self {
185 self.metadata.origin = Some(origin.into());
186 self
187 }
188
189 pub fn filename(mut self, filename: impl Into<String>) -> Self {
191 self.metadata.filename = Some(filename.into());
192 self
193 }
194
195 pub fn set_metadata(mut self, metadata: impl Into<BTreeMap<String, String>>) -> Self {
197 self.metadata.custom = metadata.into();
198 self
199 }
200
201 pub fn append_metadata(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
203 self.metadata.custom.insert(key.into(), value.into());
204 self
205 }
206
207 pub async fn send(self) -> crate::Result<Option<ResumableUpload>> {
212 let method = if self.key.is_some() {
213 Method::PUT
214 } else {
215 Method::POST
216 };
217 let request = self
218 .session
219 .request(method, self.key.as_deref().unwrap_or_default())?
220 .query(&UploadTypeQuery {
221 upload_type: UploadType::Resumable,
222 })
223 .headers(self.metadata.to_headers("")?)
224 .header(HEADER_UPLOAD_LENGTH, self.total_length.to_string());
225 let response = request.send().await?;
226
227 match response.status() {
228 StatusCode::OK => {}
229 StatusCode::NOT_IMPLEMENTED => {
230 response.drain_body().await;
231 return Ok(None);
232 }
233 status => {
234 let response = response.error_for_status_and_drain().await?;
235 response.drain_body().await;
236 return Err(Error::MalformedResponse(format!(
237 "unexpected HTTP status {status} while creating a resumable upload"
238 )));
239 }
240 }
241
242 let response: CreateSessionResponse = response.json().await?;
243 Ok(Some(
244 self.session.resume_upload(response.key, response.session),
245 ))
246 }
247}
248
249#[derive(Debug)]
251pub struct UploadProgressBuilder {
252 upload: ResumableUpload,
253}
254
255impl UploadProgressBuilder {
256 pub async fn send(self) -> crate::Result<UploadProgress> {
263 let response = self
264 .upload
265 .request(Method::PUT)?
266 .header(HEADER_UPLOAD_OFFSET, "*")
267 .send()
268 .await?;
269 parse_progress_response(response).await
270 }
271}
272
273pub struct PutChunkBuilder {
275 upload: ResumableUpload,
276 offset: u64,
277 chunk: Bytes,
278}
279
280impl fmt::Debug for PutChunkBuilder {
281 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
282 f.debug_struct("PutChunkBuilder")
283 .field("upload", &self.upload)
284 .field("offset", &self.offset)
285 .field("content_length", &self.chunk.len())
286 .finish()
287 }
288}
289
290impl PutChunkBuilder {
291 pub async fn send(self) -> crate::Result<UploadProgress> {
310 let content_length = self.chunk.len();
311 let response = self
312 .upload
313 .request(Method::PUT)?
314 .header(HEADER_UPLOAD_OFFSET, self.offset.to_string())
315 .header(reqwest::header::CONTENT_LENGTH, content_length)
316 .body(self.chunk)
317 .send()
318 .await?;
319 parse_progress_response(response).await
320 }
321}
322
323#[derive(Debug)]
325pub struct CancelUploadBuilder {
326 upload: ResumableUpload,
327}
328
329impl CancelUploadBuilder {
330 pub async fn send(self) -> crate::Result<()> {
337 let response = self.upload.request(Method::DELETE)?.send().await?;
338 match response.status() {
339 StatusCode::NO_CONTENT => {
340 response.drain_body().await;
341 Ok(())
342 }
343 StatusCode::NOT_FOUND | StatusCode::GONE => {
344 response.drain_body().await;
345 Err(Error::ResumableUploadUnavailable)
346 }
347 status => {
348 let response = response.error_for_status_and_drain().await?;
349 response.drain_body().await;
350 Err(Error::MalformedResponse(format!(
351 "unexpected HTTP status {status} while canceling a resumable upload"
352 )))
353 }
354 }
355 }
356}
357
358async fn parse_progress_response(response: Response) -> crate::Result<UploadProgress> {
359 match response.status() {
360 StatusCode::NO_CONTENT | StatusCode::CONFLICT => {
361 let offset = parse_offset(&response);
362 response.drain_body().await;
363 let offset = offset.ok_or_else(|| {
364 crate::Error::MalformedResponse(
365 "resumable upload response has no valid Upload-Offset header".into(),
366 )
367 })?;
368 Ok(UploadProgress::Incomplete { offset })
369 }
370 StatusCode::CREATED => {
371 let _: CompleteUploadResponse = response.json().await?;
372 Ok(UploadProgress::Complete)
373 }
374 StatusCode::NOT_FOUND | StatusCode::GONE => {
375 response.drain_body().await;
376 Err(Error::ResumableUploadUnavailable)
377 }
378 status => {
379 let response = response.error_for_status_and_drain().await?;
380 response.drain_body().await;
381 Err(Error::MalformedResponse(format!(
382 "unexpected HTTP status {status} while continuing a resumable upload"
383 )))
384 }
385 }
386}
387
388fn parse_offset(response: &Response) -> Option<u64> {
389 let value = response
390 .headers()
391 .get(HEADER_UPLOAD_OFFSET)?
392 .to_str()
393 .ok()?;
394 match value.parse().ok()? {
395 UploadOffset::At(offset) => Some(offset),
396 UploadOffset::Unknown => None,
397 }
398}