google_cloud_storage/storage/write_object.rs
1// Copyright 2025 Google LLC
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// https://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Contains the request builder for [write_object()] and related types.
16//!
17//! [write_object()]: crate::storage::client::Storage::write_object()
18
19use super::streaming_source::{Seek, StreamingSource};
20use super::*;
21use crate::model_ext::KeyAes256;
22use crate::storage::checksum::details::update as checksum_update;
23use crate::storage::checksum::details::{Checksum, Md5};
24use crate::storage::request_options::RequestOptions;
25
26/// A request builder for object writes.
27///
28/// # Example: hello world
29/// ```
30/// use google_cloud_storage::client::Storage;
31/// async fn sample(client: &Storage) -> anyhow::Result<()> {
32/// let response = client
33/// .write_object("projects/_/buckets/my-bucket", "hello", "Hello World!")
34/// .send_unbuffered()
35/// .await?;
36/// println!("response details={response:?}");
37/// Ok(())
38/// }
39/// ```
40///
41/// # Example: upload a file
42/// ```
43/// use google_cloud_storage::client::Storage;
44/// async fn sample(client: &Storage) -> anyhow::Result<()> {
45/// let payload = tokio::fs::File::open("my-data").await?;
46/// let response = client
47/// .write_object("projects/_/buckets/my-bucket", "my-object", payload)
48/// .send_unbuffered()
49/// .await?;
50/// println!("response details={response:?}");
51/// Ok(())
52/// }
53/// ```
54///
55/// # Example: create a new object from a custom data source
56/// ```
57/// use google_cloud_storage::{client::Storage, streaming_source::StreamingSource};
58/// struct DataSource;
59/// impl StreamingSource for DataSource {
60/// type Error = std::io::Error;
61/// async fn next(&mut self) -> Option<Result<bytes::Bytes, Self::Error>> {
62/// # panic!();
63/// }
64/// }
65///
66/// async fn sample(client: &Storage) -> anyhow::Result<()> {
67/// let response = client
68/// .write_object("projects/_/buckets/my-bucket", "my-object", DataSource)
69/// .send_buffered()
70/// .await?;
71/// println!("response details={response:?}");
72/// Ok(())
73/// }
74/// ```
75pub struct WriteObject<T, S = crate::storage::transport::Storage>
76where
77 S: crate::storage::stub::Storage + 'static,
78{
79 stub: std::sync::Arc<S>,
80 pub(crate) request: crate::model_ext::WriteObjectRequest,
81 pub(crate) payload: Payload<T>,
82 pub(crate) options: RequestOptions,
83}
84
85impl<T, S> WriteObject<T, S>
86where
87 S: crate::storage::stub::Storage + 'static,
88{
89 /// Set a [request precondition] on the object generation to match.
90 ///
91 /// With this precondition the request fails if the current object
92 /// generation matches the provided value. A common value is `0`, which
93 /// prevents writes from succeeding if the object already exists.
94 ///
95 /// # Example
96 /// ```
97 /// # use google_cloud_storage::client::Storage;
98 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
99 /// let response = client
100 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
101 /// .set_if_generation_match(0)
102 /// .send_buffered()
103 /// .await?;
104 /// println!("response details={response:?}");
105 /// # Ok(()) }
106 /// ```
107 ///
108 /// [request precondition]: https://cloud.google.com/storage/docs/request-preconditions
109 pub fn set_if_generation_match<V>(mut self, v: V) -> Self
110 where
111 V: Into<i64>,
112 {
113 self.request.spec.if_generation_match = Some(v.into());
114 self
115 }
116
117 /// Set a [request precondition] on the object generation to match.
118 ///
119 /// With this precondition the request fails if the current object
120 /// generation does not match the provided value.
121 ///
122 /// # Example
123 /// ```
124 /// # use google_cloud_storage::client::Storage;
125 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
126 /// let response = client
127 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
128 /// .set_if_generation_not_match(0)
129 /// .send_buffered()
130 /// .await?;
131 /// println!("response details={response:?}");
132 /// # Ok(()) }
133 /// ```
134 ///
135 /// [request precondition]: https://cloud.google.com/storage/docs/request-preconditions
136 pub fn set_if_generation_not_match<V>(mut self, v: V) -> Self
137 where
138 V: Into<i64>,
139 {
140 self.request.spec.if_generation_not_match = Some(v.into());
141 self
142 }
143
144 /// Set a [request precondition] on the object meta generation.
145 ///
146 /// With this precondition the request fails if the current object metadata
147 /// generation does not match the provided value. This may be useful to
148 /// prevent changes when the metageneration is known.
149 ///
150 /// # Example
151 /// ```
152 /// # use google_cloud_storage::client::Storage;
153 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
154 /// let response = client
155 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
156 /// .set_if_metageneration_match(1234)
157 /// .send_buffered()
158 /// .await?;
159 /// println!("response details={response:?}");
160 /// # Ok(()) }
161 /// ```
162 ///
163 /// [request precondition]: https://cloud.google.com/storage/docs/request-preconditions
164 pub fn set_if_metageneration_match<V>(mut self, v: V) -> Self
165 where
166 V: Into<i64>,
167 {
168 self.request.spec.if_metageneration_match = Some(v.into());
169 self
170 }
171
172 /// Set a [request precondition] on the object meta-generation.
173 ///
174 /// With this precondition the request fails if the current object metadata
175 /// generation matches the provided value. This is rarely useful in uploads,
176 /// it is more commonly used on reads to prevent a large response if the
177 /// data is already cached.
178 ///
179 /// # Example
180 /// ```
181 /// # use google_cloud_storage::client::Storage;
182 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
183 /// let response = client
184 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
185 /// .set_if_metageneration_not_match(1234)
186 /// .send_buffered()
187 /// .await?;
188 /// println!("response details={response:?}");
189 /// # Ok(()) }
190 /// ```
191 ///
192 /// [request precondition]: https://cloud.google.com/storage/docs/request-preconditions
193 pub fn set_if_metageneration_not_match<V>(mut self, v: V) -> Self
194 where
195 V: Into<i64>,
196 {
197 self.request.spec.if_metageneration_not_match = Some(v.into());
198 self
199 }
200
201 /// Sets the ACL for the new object.
202 ///
203 /// # Example
204 /// ```
205 /// # use google_cloud_storage::client::Storage;
206 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
207 /// # use google_cloud_storage::model::ObjectAccessControl;
208 /// let response = client
209 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
210 /// .set_acl([ObjectAccessControl::new().set_entity("allAuthenticatedUsers").set_role("READER")])
211 /// .send_buffered()
212 /// .await?;
213 /// println!("response details={response:?}");
214 /// # Ok(()) }
215 /// ```
216 pub fn set_acl<I, V>(mut self, v: I) -> Self
217 where
218 I: IntoIterator<Item = V>,
219 V: Into<crate::model::ObjectAccessControl>,
220 {
221 self.mut_resource().acl = v.into_iter().map(|a| a.into()).collect();
222 self
223 }
224
225 /// Sets the [cache control] for the new object.
226 ///
227 /// This can be used to control caching in [public objects].
228 ///
229 /// # Example
230 /// ```
231 /// # use google_cloud_storage::client::Storage;
232 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
233 /// let response = client
234 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
235 /// .set_cache_control("public; max-age=7200")
236 /// .send_buffered()
237 /// .await?;
238 /// println!("response details={response:?}");
239 /// # Ok(()) }
240 /// ```
241 ///
242 /// [public objects]: https://cloud.google.com/storage/docs/access-control/making-data-public
243 /// [cache control]: https://datatracker.ietf.org/doc/html/rfc7234#section-5.2
244 pub fn set_cache_control<V: Into<String>>(mut self, v: V) -> Self {
245 self.mut_resource().cache_control = v.into();
246 self
247 }
248
249 /// Sets the [content disposition] for the new object.
250 ///
251 /// Google Cloud Storage can serve content directly to web browsers. This
252 /// attribute sets the `Content-Disposition` header, which may change how
253 /// the browser displays the contents.
254 ///
255 /// # Example
256 /// ```
257 /// # use google_cloud_storage::client::Storage;
258 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
259 /// let response = client
260 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
261 /// .set_content_disposition("inline")
262 /// .send_buffered()
263 /// .await?;
264 /// println!("response details={response:?}");
265 /// # Ok(()) }
266 /// ```
267 ///
268 /// [content disposition]: https://datatracker.ietf.org/doc/html/rfc6266
269 pub fn set_content_disposition<V: Into<String>>(mut self, v: V) -> Self {
270 self.mut_resource().content_disposition = v.into();
271 self
272 }
273
274 /// Sets the [content encoding] for the object data.
275 ///
276 /// This can be used to upload compressed data and enable [transcoding] of
277 /// the data during reads.
278 ///
279 /// # Example
280 /// ```
281 /// # use google_cloud_storage::client::Storage;
282 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
283 /// use flate2::write::GzEncoder;
284 /// use std::io::Write;
285 /// let mut e = GzEncoder::new(Vec::new(), flate2::Compression::default());
286 /// e.write_all(b"hello world");
287 /// let response = client
288 /// .write_object("projects/_/buckets/my-bucket", "my-object", bytes::Bytes::from_owner(e.finish()?))
289 /// .set_content_encoding("gzip")
290 /// .send_buffered()
291 /// .await?;
292 /// println!("response details={response:?}");
293 /// # Ok(()) }
294 /// ```
295 ///
296 /// [transcoding]: https://cloud.google.com/storage/docs/transcoding
297 /// [content encoding]: https://datatracker.ietf.org/doc/html/rfc7231#section-3.1.2.2
298 pub fn set_content_encoding<V: Into<String>>(mut self, v: V) -> Self {
299 self.mut_resource().content_encoding = v.into();
300 self
301 }
302
303 /// Sets the [content language] for the new object.
304 ///
305 /// Google Cloud Storage can serve content directly to web browsers. This
306 /// attribute sets the `Content-Language` header, which may change how the
307 /// browser displays the contents.
308 ///
309 /// # Example
310 /// ```
311 /// # use google_cloud_storage::client::Storage;
312 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
313 /// let response = client
314 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
315 /// .set_content_language("en")
316 /// .send_buffered()
317 /// .await?;
318 /// println!("response details={response:?}");
319 /// # Ok(()) }
320 /// ```
321 ///
322 /// [content language]: https://cloud.google.com/storage/docs/metadata#content-language
323 pub fn set_content_language<V: Into<String>>(mut self, v: V) -> Self {
324 self.mut_resource().content_language = v.into();
325 self
326 }
327
328 /// Sets the [content type] for the new object.
329 ///
330 /// Google Cloud Storage can serve content directly to web browsers. This
331 /// attribute sets the `Content-Type` header, which may change how the
332 /// browser interprets the contents.
333 ///
334 /// # Example
335 /// ```
336 /// # use google_cloud_storage::client::Storage;
337 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
338 /// let response = client
339 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
340 /// .set_content_type("text/plain")
341 /// .send_buffered()
342 /// .await?;
343 /// println!("response details={response:?}");
344 /// # Ok(()) }
345 /// ```
346 ///
347 /// [content type]: https://datatracker.ietf.org/doc/html/rfc7231#section-3.1.1.5
348 pub fn set_content_type<V: Into<String>>(mut self, v: V) -> Self {
349 self.mut_resource().content_type = v.into();
350 self
351 }
352
353 /// Sets the [custom time] for the new object.
354 ///
355 /// This field is typically set in order to use the [DaysSinceCustomTime]
356 /// condition in Object Lifecycle Management.
357 ///
358 /// # Example
359 /// ```
360 /// # use google_cloud_storage::client::Storage;
361 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
362 /// let time = wkt::Timestamp::try_from("2025-07-07T18:30:00Z")?;
363 /// let response = client
364 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
365 /// .set_custom_time(time)
366 /// .send_buffered()
367 /// .await?;
368 /// println!("response details={response:?}");
369 /// # Ok(()) }
370 /// ```
371 ///
372 /// [DaysSinceCustomTime]: https://cloud.google.com/storage/docs/lifecycle#dayssincecustomtime
373 /// [custom time]: https://cloud.google.com/storage/docs/metadata#custom-time
374 pub fn set_custom_time<V: Into<wkt::Timestamp>>(mut self, v: V) -> Self {
375 self.mut_resource().custom_time = Some(v.into());
376 self
377 }
378
379 /// Sets the [event based hold] flag for the new object.
380 ///
381 /// This field is typically set in order to prevent objects from being
382 /// deleted or modified.
383 ///
384 /// # Example
385 /// ```
386 /// # use google_cloud_storage::client::Storage;
387 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
388 /// let response = client
389 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
390 /// .set_event_based_hold(true)
391 /// .send_buffered()
392 /// .await?;
393 /// println!("response details={response:?}");
394 /// # Ok(()) }
395 /// ```
396 ///
397 /// [event based hold]: https://cloud.google.com/storage/docs/object-holds
398 pub fn set_event_based_hold<V: Into<bool>>(mut self, v: V) -> Self {
399 self.mut_resource().event_based_hold = Some(v.into());
400 self
401 }
402
403 /// Sets the [custom metadata] for the new object.
404 ///
405 /// This field is typically set to annotate the object with
406 /// application-specific metadata.
407 ///
408 /// # Example
409 /// ```
410 /// # use google_cloud_storage::client::Storage;
411 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
412 /// let time = wkt::Timestamp::try_from("2025-07-07T18:30:00Z")?;
413 /// let response = client
414 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
415 /// .set_metadata([("test-only", "true"), ("environment", "qa")])
416 /// .send_buffered()
417 /// .await?;
418 /// println!("response details={response:?}");
419 /// # Ok(()) }
420 /// ```
421 ///
422 /// [custom metadata]: https://cloud.google.com/storage/docs/metadata#custom-metadata
423 pub fn set_metadata<I, K, V>(mut self, i: I) -> Self
424 where
425 I: IntoIterator<Item = (K, V)>,
426 K: Into<String>,
427 V: Into<String>,
428 {
429 self.mut_resource().metadata = i.into_iter().map(|(k, v)| (k.into(), v.into())).collect();
430 self
431 }
432
433 /// Sets the [retention configuration] for the new object.
434 ///
435 /// # Example
436 /// ```
437 /// # use google_cloud_storage::client::Storage;
438 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
439 /// # use google_cloud_storage::model::object::{Retention, retention};
440 /// let response = client
441 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
442 /// .set_retention(
443 /// Retention::new()
444 /// .set_mode(retention::Mode::Locked)
445 /// .set_retain_until_time(wkt::Timestamp::try_from("2035-01-01T00:00:00Z")?))
446 /// .send_buffered()
447 /// .await?;
448 /// println!("response details={response:?}");
449 /// # Ok(()) }
450 /// ```
451 ///
452 /// [retention configuration]: https://cloud.google.com/storage/docs/metadata#retention-config
453 pub fn set_retention<V>(mut self, v: V) -> Self
454 where
455 V: Into<crate::model::object::Retention>,
456 {
457 self.mut_resource().retention = Some(v.into());
458 self
459 }
460
461 /// Sets the [storage class] for the new object.
462 ///
463 /// # Example
464 /// ```
465 /// # use google_cloud_storage::client::Storage;
466 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
467 /// let response = client
468 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
469 /// .set_storage_class("ARCHIVE")
470 /// .send_buffered()
471 /// .await?;
472 /// println!("response details={response:?}");
473 /// # Ok(()) }
474 /// ```
475 ///
476 /// [storage class]: https://cloud.google.com/storage/docs/storage-classes
477 pub fn set_storage_class<V>(mut self, v: V) -> Self
478 where
479 V: Into<String>,
480 {
481 self.mut_resource().storage_class = v.into();
482 self
483 }
484
485 /// Sets the [temporary hold] flag for the new object.
486 ///
487 /// This field is typically set in order to prevent objects from being
488 /// deleted or modified.
489 ///
490 /// # Example
491 /// ```
492 /// # use google_cloud_storage::client::Storage;
493 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
494 /// let time = wkt::Timestamp::try_from("2025-07-07T18:30:00Z")?;
495 /// let response = client
496 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
497 /// .set_temporary_hold(true)
498 /// .send_buffered()
499 /// .await?;
500 /// println!("response details={response:?}");
501 /// # Ok(()) }
502 /// ```
503 ///
504 /// [temporary hold]: https://cloud.google.com/storage/docs/object-holds
505 pub fn set_temporary_hold<V: Into<bool>>(mut self, v: V) -> Self {
506 self.mut_resource().temporary_hold = v.into();
507 self
508 }
509
510 /// Sets the resource name of the [Customer-managed encryption key] for this
511 /// object.
512 ///
513 /// The service imposes a number of restrictions on the keys used to encrypt
514 /// Google Cloud Storage objects. Read the documentation in full before
515 /// trying to use customer-managed encryption keys. In particular, verify
516 /// the service has the necessary permissions, and the key is in a
517 /// compatible location.
518 ///
519 /// # Example
520 /// ```
521 /// # use google_cloud_storage::client::Storage;
522 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
523 /// let response = client
524 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
525 /// .set_kms_key("projects/test-project/locations/us-central1/keyRings/test-ring/cryptoKeys/test-key")
526 /// .send_buffered()
527 /// .await?;
528 /// println!("response details={response:?}");
529 /// # Ok(()) }
530 /// ```
531 ///
532 /// [Customer-managed encryption key]: https://cloud.google.com/storage/docs/encryption/customer-managed-keys
533 pub fn set_kms_key<V>(mut self, v: V) -> Self
534 where
535 V: Into<String>,
536 {
537 self.mut_resource().kms_key = v.into();
538 self
539 }
540
541 /// Configure this object to use one of the [predefined ACLs].
542 ///
543 /// # Example
544 /// ```
545 /// # use google_cloud_storage::client::Storage;
546 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
547 /// let response = client
548 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
549 /// .set_predefined_acl("private")
550 /// .send_buffered()
551 /// .await?;
552 /// println!("response details={response:?}");
553 /// # Ok(()) }
554 /// ```
555 ///
556 /// [predefined ACLs]: https://cloud.google.com/storage/docs/access-control/lists#predefined-acl
557 pub fn set_predefined_acl<V>(mut self, v: V) -> Self
558 where
559 V: Into<String>,
560 {
561 self.request.spec.predefined_acl = v.into();
562 self
563 }
564
565 /// The encryption key used with the Customer-Supplied Encryption Keys
566 /// feature. In raw bytes format (not base64-encoded).
567 ///
568 /// # Example
569 /// ```
570 /// # use google_cloud_storage::client::Storage;
571 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
572 /// # use google_cloud_storage::model_ext::KeyAes256;
573 /// let key: &[u8] = &[97; 32];
574 /// let response = client
575 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
576 /// .set_key(KeyAes256::new(key)?)
577 /// .send_buffered()
578 /// .await?;
579 /// println!("response details={response:?}");
580 /// # Ok(()) }
581 /// ```
582 pub fn set_key(mut self, v: KeyAes256) -> Self {
583 self.request.params = Some(v.into());
584 self
585 }
586
587 /// Sets the object custom contexts.
588 ///
589 /// # Example
590 /// ```
591 /// # use google_cloud_storage::client::Storage;
592 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
593 /// # use google_cloud_storage::model::{ObjectContexts, ObjectCustomContextPayload};
594 /// let response = client
595 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
596 /// .set_contexts(
597 /// ObjectContexts::new().set_custom([
598 /// ("example", ObjectCustomContextPayload::new().set_value("true")),
599 /// ])
600 /// )
601 /// .send_buffered()
602 /// .await?;
603 /// println!("response details={response:?}");
604 /// # Ok(()) }
605 /// ```
606 pub fn set_contexts<V>(mut self, v: V) -> Self
607 where
608 V: Into<crate::model::ObjectContexts>,
609 {
610 self.mut_resource().contexts = Some(v.into());
611 self
612 }
613
614 /// Configure the idempotency for this upload.
615 ///
616 /// By default, the client library treats single-shot uploads without
617 /// preconditions, as non-idempotent. If the destination bucket is
618 /// configured with [object versioning] then the operation may succeed
619 /// multiple times with observable side-effects. With object versioning and
620 /// a [lifecycle] policy limiting the number of versions, uploading the same
621 /// data multiple times may result in data loss.
622 ///
623 /// The client library cannot efficiently determine if these conditions
624 /// apply to your upload. If they do, or your application can tolerate
625 /// multiple versions of the same data for other reasons, consider using
626 /// `with_idempotency(true)`.
627 ///
628 /// The client library treats resumable uploads as idempotent, regardless of
629 /// the value in this option. Such uploads can succeed at most once.
630 ///
631 /// # Example
632 /// ```
633 /// # use google_cloud_storage::client::Storage;
634 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
635 /// use std::time::Duration;
636 /// use google_cloud_gax::retry_policy::RetryPolicyExt;
637 /// let response = client
638 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
639 /// .with_idempotency(true)
640 /// .send_buffered()
641 /// .await?;
642 /// println!("response details={response:?}");
643 /// # Ok(()) }
644 /// ```
645 ///
646 /// [lifecycle]: https://cloud.google.com/storage/docs/lifecycle
647 /// [object versioning]: https://cloud.google.com/storage/docs/object-versioning
648 pub fn with_idempotency(mut self, v: bool) -> Self {
649 self.options.idempotency = Some(v);
650 self
651 }
652
653 /// The retry policy used for this request.
654 ///
655 /// # Example
656 /// ```
657 /// # use google_cloud_storage::client::Storage;
658 /// # use google_cloud_storage::retry_policy::RetryableErrors;
659 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
660 /// use std::time::Duration;
661 /// use google_cloud_gax::retry_policy::RetryPolicyExt;
662 /// let response = client
663 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
664 /// .with_retry_policy(
665 /// RetryableErrors
666 /// .with_attempt_limit(5)
667 /// .with_time_limit(Duration::from_secs(90)),
668 /// )
669 /// .send_buffered()
670 /// .await?;
671 /// println!("response details={response:?}");
672 /// # Ok(()) }
673 /// ```
674 pub fn with_retry_policy<V: Into<google_cloud_gax::retry_policy::RetryPolicyArg>>(
675 mut self,
676 v: V,
677 ) -> Self {
678 self.options.retry_policy = v.into().into();
679 self
680 }
681
682 /// The backoff policy used for this request.
683 ///
684 /// # Example
685 /// ```
686 /// # use google_cloud_storage::client::Storage;
687 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
688 /// use std::time::Duration;
689 /// use google_cloud_gax::exponential_backoff::ExponentialBackoff;
690 /// let response = client
691 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
692 /// .with_backoff_policy(ExponentialBackoff::default())
693 /// .send_buffered()
694 /// .await?;
695 /// println!("response details={response:?}");
696 /// # Ok(()) }
697 /// ```
698 pub fn with_backoff_policy<V: Into<google_cloud_gax::backoff_policy::BackoffPolicyArg>>(
699 mut self,
700 v: V,
701 ) -> Self {
702 self.options.backoff_policy = v.into().into();
703 self
704 }
705
706 /// The retry throttler used for this request.
707 ///
708 /// Most of the time you want to use the same throttler for all the requests
709 /// in a client, and even the same throttler for many clients. Rarely it
710 /// may be necessary to use an custom throttler for some subset of the
711 /// requests.
712 ///
713 /// # Example
714 /// ```
715 /// # use google_cloud_storage::client::Storage;
716 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
717 /// let response = client
718 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
719 /// .with_retry_throttler(adhoc_throttler())
720 /// .send_buffered()
721 /// .await?;
722 /// println!("response details={response:?}");
723 /// fn adhoc_throttler() -> google_cloud_gax::retry_throttler::SharedRetryThrottler {
724 /// # panic!();
725 /// }
726 /// # Ok(()) }
727 /// ```
728 pub fn with_retry_throttler<V: Into<google_cloud_gax::retry_throttler::RetryThrottlerArg>>(
729 mut self,
730 v: V,
731 ) -> Self {
732 self.options.retry_throttler = v.into().into();
733 self
734 }
735
736 /// Sets the payload size threshold to switch from single-shot to resumable uploads.
737 ///
738 /// # Example
739 /// ```
740 /// # use google_cloud_storage::client::Storage;
741 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
742 /// let response = client
743 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
744 /// .with_resumable_upload_threshold(0_usize) // Forces a resumable upload.
745 /// .send_buffered()
746 /// .await?;
747 /// println!("response details={response:?}");
748 /// # Ok(()) }
749 /// ```
750 ///
751 /// The client library can perform uploads using [single-shot] or
752 /// [resumable] uploads. For small objects, single-shot uploads offer better
753 /// performance, as they require a single HTTP transfer. For larger objects,
754 /// the additional request latency is not significant, and resumable uploads
755 /// offer better recovery on errors.
756 ///
757 /// The library automatically selects resumable uploads when the payload is
758 /// equal to or larger than this option. For smaller uploads the client
759 /// library uses single-shot uploads.
760 ///
761 /// The exact threshold depends on where the application is deployed and
762 /// destination bucket location with respect to where the application is
763 /// running. The library defaults should work well in most cases, but some
764 /// applications may benefit from fine-tuning.
765 ///
766 /// [single-shot]: https://cloud.google.com/storage/docs/uploading-objects
767 /// [resumable]: https://cloud.google.com/storage/docs/resumable-uploads
768 pub fn with_resumable_upload_threshold<V: Into<usize>>(mut self, v: V) -> Self {
769 self.options.set_resumable_upload_threshold(v.into());
770 self
771 }
772
773 /// Changes the buffer size for some resumable uploads.
774 ///
775 /// # Example
776 /// ```
777 /// # use google_cloud_storage::client::Storage;
778 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
779 /// let response = client
780 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
781 /// .with_resumable_upload_buffer_size(32 * 1024 * 1024_usize)
782 /// .send_buffered()
783 /// .await?;
784 /// println!("response details={response:?}");
785 /// # Ok(()) }
786 /// ```
787 ///
788 /// When performing [resumable uploads] from sources without [Seek] the
789 /// client library needs to buffer data in memory until it is persisted by
790 /// the service. Otherwise the data would be lost if the upload fails.
791 /// Applications may want to tune this buffer size:
792 ///
793 /// - Use smaller buffer sizes to support more concurrent uploads in the
794 /// same application.
795 /// - Use larger buffer sizes for better throughput. Sending many small
796 /// buffers stalls the upload until the client receives a successful
797 /// response from the service.
798 ///
799 /// Keep in mind that there are diminishing returns on using larger buffers.
800 ///
801 /// [resumable uploads]: https://cloud.google.com/storage/docs/resumable-uploads
802 /// [Seek]: crate::streaming_source::Seek
803 pub fn with_resumable_upload_buffer_size<V: Into<usize>>(mut self, v: V) -> Self {
804 self.options.set_resumable_upload_buffer_size(v.into());
805 self
806 }
807
808 /// Sets the `User-Agent` header for this request.
809 ///
810 /// # Example
811 /// ```
812 /// # use google_cloud_storage::client::Storage;
813 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
814 /// let mut response = client
815 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
816 /// .with_user_agent("my-app/1.0.0")
817 /// .send_buffered()
818 /// .await?;
819 /// println!("response details={response:?}");
820 /// # Ok(()) }
821 /// ```
822 pub fn with_user_agent(mut self, user_agent: impl Into<String>) -> Self {
823 self.options.user_agent = Some(user_agent.into());
824 self
825 }
826
827 /// Sets the project that will be billed for this request.
828 ///
829 /// Required for [Requester Pays] buckets. The value overrides any
830 /// `quota_project_id` configured on the credentials; the credential-level
831 /// header is suppressed for this RPC.
832 ///
833 /// # Example
834 /// ```
835 /// # use google_cloud_storage::client::Storage;
836 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
837 /// let response = client
838 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello")
839 /// .with_quota_project("my-billing-project")
840 /// .send_buffered()
841 /// .await?;
842 /// # Ok(()) }
843 /// ```
844 ///
845 /// [Requester Pays]: https://cloud.google.com/storage/docs/requester-pays
846 pub fn with_quota_project(mut self, project: impl Into<String>) -> Self {
847 self.options.set_quota_project(project);
848 self
849 }
850
851 fn mut_resource(&mut self) -> &mut crate::model::Object {
852 self.request
853 .spec
854 .resource
855 .as_mut()
856 .expect("resource field initialized in `new()`")
857 }
858
859 fn set_crc32c<V: Into<u32>>(mut self, v: V) -> Self {
860 let checksum = self.mut_resource().checksums.get_or_insert_default();
861 checksum.crc32c = Some(v.into());
862 self
863 }
864
865 /// Sets the MD5 hash for the object being written.
866 pub fn set_md5_hash<I, V>(mut self, i: I) -> Self
867 where
868 I: IntoIterator<Item = V>,
869 V: Into<u8>,
870 {
871 let checksum = self.mut_resource().checksums.get_or_insert_default();
872 checksum.md5_hash = i.into_iter().map(|v| v.into()).collect();
873 self
874 }
875
876 /// Provide a precomputed value for the CRC32C checksum.
877 ///
878 /// # Example
879 /// ```
880 /// # use google_cloud_storage::client::Storage;
881 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
882 /// use crc32c::crc32c;
883 /// let response = client
884 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
885 /// .with_known_crc32c(crc32c(b"hello world"))
886 /// .send_buffered()
887 /// .await?;
888 /// println!("response details={response:?}");
889 /// # Ok(()) }
890 /// ```
891 ///
892 /// In some applications, the payload's CRC32C checksum is already known.
893 /// For example, the application may be reading the data from another blob
894 /// storage system.
895 ///
896 /// In such cases, it is safer to pass the known CRC32C of the payload to
897 /// [Cloud Storage], and more efficient to skip the computation in the
898 /// client library.
899 ///
900 /// Note that once you provide a CRC32C value to this builder you cannot
901 /// use [compute_md5()] to also have the library compute the checksums.
902 ///
903 /// [compute_md5()]: WriteObject::compute_md5
904 pub fn with_known_crc32c<V: Into<u32>>(self, v: V) -> Self {
905 let mut this = self;
906 this.options.checksum.crc32c = None;
907 this.set_crc32c(v)
908 }
909
910 /// Provide a precomputed value for the MD5 hash.
911 ///
912 /// # Example
913 /// ```
914 /// # use google_cloud_storage::client::Storage;
915 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
916 /// use md5::compute;
917 /// let hash = md5::compute(b"hello world");
918 /// let response = client
919 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
920 /// .with_known_md5_hash(bytes::Bytes::from_owner(hash.0))
921 /// .send_buffered()
922 /// .await?;
923 /// println!("response details={response:?}");
924 /// # Ok(()) }
925 /// ```
926 ///
927 /// In some applications, the payload's MD5 hash is already known. For
928 /// example, the application may be reading the data from another blob
929 /// storage system.
930 ///
931 /// In such cases, it is safer to pass the known MD5 of the payload to
932 /// [Cloud Storage], and more efficient to skip the computation in the
933 /// client library.
934 ///
935 /// Note that once you provide a MD5 value to this builder you cannot
936 /// use [compute_md5()] to also have the library compute the checksums.
937 ///
938 /// [compute_md5()]: WriteObject::compute_md5
939 pub fn with_known_md5_hash<I, V>(self, i: I) -> Self
940 where
941 I: IntoIterator<Item = V>,
942 V: Into<u8>,
943 {
944 let mut this = self;
945 this.options.checksum.md5_hash = None;
946 this.set_md5_hash(i)
947 }
948
949 /// Enables computation of MD5 hashes.
950 ///
951 /// # Example
952 /// ```
953 /// # use google_cloud_storage::client::Storage;
954 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
955 /// let payload = tokio::fs::File::open("my-data").await?;
956 /// let response = client
957 /// .write_object("projects/_/buckets/my-bucket", "my-object", payload)
958 /// .compute_md5()
959 /// .send_buffered()
960 /// .await?;
961 /// println!("response details={response:?}");
962 /// # Ok(()) }
963 /// ```
964 ///
965 /// See [precompute_checksums][WriteObject::precompute_checksums] for more
966 /// details on how checksums are used by the client library and their
967 /// limitations.
968 pub fn compute_md5(self) -> Self {
969 let mut this = self;
970 this.options.checksum.md5_hash = Some(Md5::default());
971 this
972 }
973
974 pub(crate) fn new<B, O, P>(
975 stub: std::sync::Arc<S>,
976 bucket: B,
977 object: O,
978 payload: P,
979 options: RequestOptions,
980 ) -> Self
981 where
982 B: Into<String>,
983 O: Into<String>,
984 P: Into<Payload<T>>,
985 {
986 let resource = crate::model::Object::new()
987 .set_bucket(bucket)
988 .set_name(object);
989 WriteObject {
990 stub,
991 request: crate::model_ext::WriteObjectRequest {
992 spec: crate::model::WriteObjectSpec::new().set_resource(resource),
993 params: None,
994 checksum_precomputation: true,
995 },
996 payload: payload.into(),
997 options,
998 }
999 }
1000}
1001
1002impl<T, S> WriteObject<T, S>
1003where
1004 T: StreamingSource + Seek + Send + Sync + 'static,
1005 <T as StreamingSource>::Error: std::error::Error + Send + Sync + 'static,
1006 <T as Seek>::Error: std::error::Error + Send + Sync + 'static,
1007 S: crate::storage::stub::Storage + 'static,
1008{
1009 /// A simple upload from a buffer.
1010 ///
1011 /// # Example
1012 /// ```
1013 /// # use google_cloud_storage::client::Storage;
1014 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
1015 /// let response = client
1016 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
1017 /// .send_unbuffered()
1018 /// .await?;
1019 /// println!("response details={response:?}");
1020 /// # Ok(()) }
1021 /// ```
1022 pub async fn send_unbuffered(self) -> Result<Object> {
1023 self.stub
1024 .write_object_unbuffered(self.payload, self.request, self.options)
1025 .await
1026 }
1027
1028 /// Configures checksum precomputation for unbuffered resumable uploads.
1029 ///
1030 /// Checksum precomputation is **enabled by default (`true`)** to ensure server-side data
1031 /// integrity validation for large objects streamed via resumable uploads.
1032 ///
1033 /// Note: Single-shot unbuffered uploads (< resumable upload threshold) always compute checksums
1034 /// on the fly via trailing multipart metadata and do not require precomputation.
1035 ///
1036 /// Call `.with_checksum_precomputation(false)` to turn off precomputation and prioritize
1037 /// upload speed.
1038 ///
1039 /// # Example
1040 /// ```
1041 /// # use google_cloud_storage::client::Storage;
1042 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
1043 /// let payload = tokio::fs::File::open("my-data").await?;
1044 /// let response = client
1045 /// .write_object("projects/_/buckets/my-bucket", "my-object", payload)
1046 /// .with_checksum_precomputation(false)
1047 /// .send_unbuffered()
1048 /// .await?;
1049 /// println!("response details={response:?}");
1050 /// # Ok(()) }
1051 /// ```
1052 pub fn with_checksum_precomputation(mut self, enable: bool) -> Self {
1053 self.request.checksum_precomputation = enable;
1054 self
1055 }
1056
1057 /// Precomputes the payload checksums before uploading the data.
1058 ///
1059 /// If the checksums are known when the upload starts, the client library
1060 /// includes them with the upload request, allowing the service to reject
1061 /// the upload if the payload and the checksums do not match.
1062 ///
1063 /// # Example
1064 /// ```
1065 /// # use google_cloud_storage::client::Storage;
1066 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
1067 /// let payload = tokio::fs::File::open("my-data").await?;
1068 /// #[allow(deprecated)]
1069 /// let response = client
1070 /// .write_object("projects/_/buckets/my-bucket", "my-object", payload)
1071 /// .precompute_checksums()
1072 /// .await?
1073 /// .send_unbuffered()
1074 /// .await?;
1075 /// println!("response details={response:?}");
1076 /// # Ok(()) }
1077 /// ```
1078 ///
1079 /// # Deprecated
1080 /// `precompute_checksums()` is now redundant because checksums are automatically
1081 /// validated by default across all upload modes:
1082 /// - **Buffered upload**: Checksum is uploaded in the final chunk.
1083 /// - **Single-shot unbuffered upload**: Checksum is attached as trailing metadata in a
1084 /// multipart request.
1085 /// - **Resumable unbuffered upload**: Checksum is precomputed upfront by default
1086 /// (or supplied via `with_known_crc32c`).
1087 ///
1088 /// For unbuffered resumable uploads, precomputing checksums incurs an extra read pass.
1089 /// If you want to prioritize upload speed over data integrity, call
1090 /// [with_checksum_precomputation(false)][WriteObject::with_checksum_precomputation].
1091 #[deprecated(
1092 since = "1.19.0",
1093 note = "`precompute_checksums()` is redundant as checksums are validated by default. \
1094 For unbuffered resumable uploads, call `with_checksum_precomputation(false)` \
1095 to prioritize speed over upfront integrity validation."
1096 )]
1097 pub async fn precompute_checksums(mut self) -> Result<Self> {
1098 let mut offset = 0_u64;
1099 self.payload.seek(offset).await.map_err(Error::ser)?;
1100 while let Some(n) = self.payload.next().await.transpose().map_err(Error::ser)? {
1101 self.options.checksum.update(offset, &n);
1102 offset += n.len() as u64;
1103 }
1104 self.payload.seek(0_u64).await.map_err(Error::ser)?;
1105 let computed = self.options.checksum.finalize();
1106 let current = self.mut_resource().checksums.get_or_insert_default();
1107 checksum_update(current, computed);
1108 self.options.checksum = Checksum {
1109 crc32c: None,
1110 md5_hash: None,
1111 };
1112 Ok(self)
1113 }
1114}
1115
1116impl<T, S> WriteObject<T, S>
1117where
1118 T: StreamingSource + Send + Sync + 'static,
1119 T::Error: std::error::Error + Send + Sync + 'static,
1120 S: crate::storage::stub::Storage + 'static,
1121{
1122 /// Upload an object from a streaming source without rewinds.
1123 ///
1124 /// If the data source does **not** implement [Seek] the client library must
1125 /// buffer data sent to the service until the service confirms it has
1126 /// persisted the data. This requires more memory in the client, and when
1127 /// the buffer grows too large, may require stalling the writer until the
1128 /// service can persist the data.
1129 ///
1130 /// Use this function for data sources where it is expensive or impossible
1131 /// to restart the data source. This function is also useful when it is hard
1132 /// or impossible to predict the number of bytes emitted by a stream, even
1133 /// if restarting the stream is not too expensive.
1134 ///
1135 /// # Example
1136 /// ```
1137 /// # use google_cloud_storage::client::Storage;
1138 /// # async fn sample(client: &Storage) -> anyhow::Result<()> {
1139 /// let response = client
1140 /// .write_object("projects/_/buckets/my-bucket", "my-object", "hello world")
1141 /// .send_buffered()
1142 /// .await?;
1143 /// println!("response details={response:?}");
1144 /// # Ok(()) }
1145 /// ```
1146 pub async fn send_buffered(self) -> crate::Result<Object> {
1147 self.stub
1148 .write_object_buffered(self.payload, self.request, self.options)
1149 .await
1150 }
1151}
1152
1153// We need `Debug` to use `expect_err()` in `Result<WriteObject, ...>`.
1154impl<T, S> std::fmt::Debug for WriteObject<T, S>
1155where
1156 S: crate::storage::stub::Storage + 'static,
1157{
1158 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1159 f.debug_struct("WriteObject")
1160 .field("stub", &self.stub)
1161 .field("request", &self.request)
1162 // skip payload, as it is not `Debug`
1163 .field("options", &self.options)
1164 .finish()
1165 }
1166}
1167
1168#[cfg(test)]
1169#[allow(deprecated)]
1170mod tests {
1171 use super::client::tests::{test_builder, test_inner_client};
1172 use super::*;
1173 use crate::client::Storage;
1174 use crate::model::{
1175 CommonObjectRequestParams, ObjectChecksums, ObjectContexts, ObjectCustomContextPayload,
1176 WriteObjectSpec,
1177 };
1178 use crate::storage::checksum::details::{Crc32c, Md5};
1179 use crate::streaming_source::tests::MockSeekSource;
1180 use google_cloud_auth::credentials::anonymous::Builder as Anonymous;
1181 use httptest::{Expectation, Server, matchers::*, responders::status_code};
1182 use std::error::Error as _;
1183 use std::io::{Error as IoError, ErrorKind};
1184
1185 type Result = anyhow::Result<()>;
1186
1187 // Verify `write_object()` can be used with a source that implements
1188 // `StreamingSource` **and** `Seek`
1189 #[tokio::test]
1190 async fn test_upload_streaming_source_and_seek() -> Result {
1191 struct Source;
1192 impl crate::streaming_source::StreamingSource for Source {
1193 type Error = std::io::Error;
1194 async fn next(&mut self) -> Option<std::result::Result<bytes::Bytes, Self::Error>> {
1195 None
1196 }
1197 }
1198 impl crate::streaming_source::Seek for Source {
1199 type Error = std::io::Error;
1200 async fn seek(&mut self, _offset: u64) -> std::result::Result<(), Self::Error> {
1201 Ok(())
1202 }
1203 }
1204
1205 let client = Storage::builder()
1206 .with_credentials(Anonymous::new().build())
1207 .build()
1208 .await?;
1209 let _ = client.write_object("projects/_/buckets/test-bucket", "test-object", Source);
1210 Ok(())
1211 }
1212
1213 // Verify `write_object()` can be used with a source that **only**
1214 // implements `StreamingSource`.
1215 #[tokio::test]
1216 async fn test_upload_only_streaming_source() -> Result {
1217 struct Source;
1218 impl crate::streaming_source::StreamingSource for Source {
1219 type Error = std::io::Error;
1220 async fn next(&mut self) -> Option<std::result::Result<bytes::Bytes, Self::Error>> {
1221 None
1222 }
1223 }
1224
1225 let client = Storage::builder()
1226 .with_credentials(Anonymous::new().build())
1227 .build()
1228 .await?;
1229 let _ = client.write_object("projects/_/buckets/test-bucket", "test-object", Source);
1230 Ok(())
1231 }
1232
1233 // Verify `write_object()` meets normal Send, Sync, requirements.
1234 #[tokio::test]
1235 async fn test_upload_is_send_and_static() -> Result {
1236 let client = Storage::builder()
1237 .with_credentials(Anonymous::new().build())
1238 .build()
1239 .await?;
1240
1241 fn need_send<T: Send>(_val: &T) {}
1242 fn need_sync<T: Sync>(_val: &T) {}
1243 fn need_static<T: 'static>(_val: &T) {}
1244
1245 let upload = client.write_object("projects/_/buckets/test-bucket", "test-object", "");
1246 need_send(&upload);
1247 need_sync(&upload);
1248 need_static(&upload);
1249
1250 let upload = client
1251 .write_object("projects/_/buckets/test-bucket", "test-object", "")
1252 .send_unbuffered();
1253 need_send(&upload);
1254 need_static(&upload);
1255
1256 let upload = client
1257 .write_object("projects/_/buckets/test-bucket", "test-object", "")
1258 .send_buffered();
1259 need_send(&upload);
1260 need_static(&upload);
1261
1262 Ok(())
1263 }
1264
1265 #[tokio::test]
1266 async fn write_object_metadata() -> Result {
1267 use crate::model::ObjectAccessControl;
1268 let inner = test_inner_client(test_builder()).await;
1269 let options = inner.options.clone();
1270 let stub = crate::storage::transport::Storage::new_test(inner);
1271 let key = KeyAes256::new(&[0x42; 32]).expect("hard-coded key is not an error");
1272 let mut builder =
1273 WriteObject::new(stub, "projects/_/buckets/bucket", "object", "", options)
1274 .set_if_generation_match(10)
1275 .set_if_generation_not_match(20)
1276 .set_if_metageneration_match(30)
1277 .set_if_metageneration_not_match(40)
1278 .set_predefined_acl("private")
1279 .set_acl([ObjectAccessControl::new()
1280 .set_entity("allAuthenticatedUsers")
1281 .set_role("READER")])
1282 .set_cache_control("public; max-age=7200")
1283 .set_content_disposition("inline")
1284 .set_content_encoding("gzip")
1285 .set_content_language("en")
1286 .set_content_type("text/plain")
1287 .set_contexts(ObjectContexts::new().set_custom([(
1288 "context-key",
1289 ObjectCustomContextPayload::new().set_value("context-value"),
1290 )]))
1291 .set_custom_time(wkt::Timestamp::try_from("2025-07-07T18:11:00Z")?)
1292 .set_event_based_hold(true)
1293 .set_key(key.clone())
1294 .set_metadata([("k0", "v0"), ("k1", "v1")])
1295 .set_retention(
1296 crate::model::object::Retention::new()
1297 .set_mode(crate::model::object::retention::Mode::Locked)
1298 .set_retain_until_time(wkt::Timestamp::try_from("2035-07-07T18:14:00Z")?),
1299 )
1300 .set_storage_class("ARCHIVE")
1301 .set_temporary_hold(true)
1302 .set_kms_key("test-key")
1303 .with_known_crc32c(crc32c::crc32c(b""))
1304 .with_known_md5_hash(md5::compute(b"").0);
1305
1306 let resource = builder.request.spec.resource.take().unwrap();
1307 let builder = builder;
1308 assert_eq!(
1309 &builder.request.spec,
1310 &WriteObjectSpec::new()
1311 .set_if_generation_match(10)
1312 .set_if_generation_not_match(20)
1313 .set_if_metageneration_match(30)
1314 .set_if_metageneration_not_match(40)
1315 .set_predefined_acl("private")
1316 );
1317
1318 assert_eq!(
1319 &builder.request.params,
1320 &Some(CommonObjectRequestParams::from(key))
1321 );
1322
1323 assert_eq!(
1324 resource,
1325 Object::new()
1326 .set_name("object")
1327 .set_bucket("projects/_/buckets/bucket")
1328 .set_acl([ObjectAccessControl::new()
1329 .set_entity("allAuthenticatedUsers")
1330 .set_role("READER")])
1331 .set_cache_control("public; max-age=7200")
1332 .set_content_disposition("inline")
1333 .set_content_encoding("gzip")
1334 .set_content_language("en")
1335 .set_content_type("text/plain")
1336 .set_contexts(ObjectContexts::new().set_custom([(
1337 "context-key",
1338 ObjectCustomContextPayload::new().set_value("context-value"),
1339 )]))
1340 .set_checksums(
1341 crate::model::ObjectChecksums::new()
1342 .set_crc32c(crc32c::crc32c(b""))
1343 .set_md5_hash(bytes::Bytes::from_iter(md5::compute(b"").0))
1344 )
1345 .set_custom_time(wkt::Timestamp::try_from("2025-07-07T18:11:00Z")?)
1346 .set_event_based_hold(true)
1347 .set_metadata([("k0", "v0"), ("k1", "v1")])
1348 .set_retention(
1349 crate::model::object::Retention::new()
1350 .set_mode("LOCKED")
1351 .set_retain_until_time(wkt::Timestamp::try_from("2035-07-07T18:14:00Z")?)
1352 )
1353 .set_storage_class("ARCHIVE")
1354 .set_temporary_hold(true)
1355 .set_kms_key("test-key")
1356 );
1357
1358 Ok(())
1359 }
1360
1361 #[tokio::test]
1362 async fn upload_object_options() {
1363 let inner = test_inner_client(
1364 test_builder()
1365 .with_resumable_upload_threshold(123_usize)
1366 .with_resumable_upload_buffer_size(234_usize),
1367 )
1368 .await;
1369 let options = inner.options.clone();
1370 let stub = crate::storage::transport::Storage::new_test(inner);
1371 let request = WriteObject::new(
1372 stub.clone(),
1373 "projects/_/buckets/bucket",
1374 "object",
1375 "",
1376 options.clone(),
1377 );
1378 assert_eq!(request.options.resumable_upload_threshold(), 123);
1379 assert_eq!(request.options.resumable_upload_buffer_size(), 234);
1380 assert_eq!(request.options.user_agent, None);
1381
1382 let user_agent = "quick_foxes_lazy_dogs/1.0.0";
1383 let request = WriteObject::new(stub, "projects/_/buckets/bucket", "object", "", options)
1384 .with_resumable_upload_threshold(345_usize)
1385 .with_resumable_upload_buffer_size(456_usize)
1386 .with_user_agent(user_agent);
1387 assert_eq!(request.options.resumable_upload_threshold(), 345);
1388 assert_eq!(request.options.resumable_upload_buffer_size(), 456);
1389 assert_eq!(request.options.user_agent.as_deref(), Some(user_agent));
1390 }
1391
1392 const QUICK: &str = "the quick brown fox jumps over the lazy dog";
1393 const VEXING: &str = "how vexingly quick daft zebras jump";
1394
1395 fn quick_checksum(mut engine: Checksum) -> ObjectChecksums {
1396 engine.update(0, &bytes::Bytes::from_static(QUICK.as_bytes()));
1397 engine.finalize()
1398 }
1399
1400 async fn collect<S: StreamingSource>(mut stream: S) -> anyhow::Result<Vec<u8>> {
1401 let mut collected = Vec::new();
1402 while let Some(b) = stream.next().await.transpose()? {
1403 collected.extend_from_slice(&b);
1404 }
1405 Ok(collected)
1406 }
1407
1408 #[tokio::test]
1409 async fn checksum_default() -> Result {
1410 let client = test_builder().build().await?;
1411 let upload = client
1412 .write_object("my-bucket", "my-object", QUICK)
1413 .precompute_checksums()
1414 .await?;
1415 let want = quick_checksum(Checksum {
1416 crc32c: Some(Crc32c::default()),
1417 md5_hash: None,
1418 });
1419 assert_eq!(
1420 upload.request.spec.resource.and_then(|r| r.checksums),
1421 Some(want)
1422 );
1423 let collected = collect(upload.payload).await?;
1424 assert_eq!(collected, QUICK.as_bytes());
1425 Ok(())
1426 }
1427
1428 #[tokio::test]
1429 async fn checksum_md5_and_crc32c() -> Result {
1430 let client = test_builder().build().await?;
1431 let upload = client
1432 .write_object("my-bucket", "my-object", QUICK)
1433 .compute_md5()
1434 .precompute_checksums()
1435 .await?;
1436 let want = quick_checksum(Checksum {
1437 crc32c: Some(Crc32c::default()),
1438 md5_hash: Some(Md5::default()),
1439 });
1440 assert_eq!(
1441 upload.request.spec.resource.and_then(|r| r.checksums),
1442 Some(want)
1443 );
1444 Ok(())
1445 }
1446
1447 #[tokio::test]
1448 async fn checksum_precomputed() -> Result {
1449 let mut engine = Checksum {
1450 crc32c: Some(Crc32c::default()),
1451 md5_hash: Some(Md5::default()),
1452 };
1453 engine.update(0, &bytes::Bytes::from_static(VEXING.as_bytes()));
1454 let ck = engine.finalize();
1455
1456 let client = test_builder().build().await?;
1457 let upload = client
1458 .write_object("my-bucket", "my-object", QUICK)
1459 .with_known_crc32c(ck.crc32c.unwrap())
1460 .with_known_md5_hash(ck.md5_hash.clone())
1461 .precompute_checksums()
1462 .await?;
1463 // Note that the checksums do not match the data. This is intentional,
1464 // we are trying to verify that whatever is provided in with_crc32c()
1465 // and with_md5() is respected.
1466 assert_eq!(
1467 upload.request.spec.resource.and_then(|r| r.checksums),
1468 Some(ck)
1469 );
1470
1471 Ok(())
1472 }
1473
1474 #[tokio::test]
1475 async fn checksum_crc32c_known_md5_computed() -> Result {
1476 let mut engine = Checksum {
1477 crc32c: Some(Crc32c::default()),
1478 md5_hash: Some(Md5::default()),
1479 };
1480 engine.update(0, &bytes::Bytes::from_static(VEXING.as_bytes()));
1481 let ck = engine.finalize();
1482
1483 let client = test_builder().build().await?;
1484 let upload = client
1485 .write_object("my-bucket", "my-object", QUICK)
1486 .compute_md5()
1487 .with_known_crc32c(ck.crc32c.unwrap())
1488 .precompute_checksums()
1489 .await?;
1490 // Note that the checksums do not match the data. This is intentional,
1491 // we are trying to verify that whatever is provided in with_known*()
1492 // is respected.
1493 let want = quick_checksum(Checksum {
1494 crc32c: None,
1495 md5_hash: Some(Md5::default()),
1496 })
1497 .set_crc32c(ck.crc32c.unwrap());
1498 assert_eq!(
1499 upload.request.spec.resource.and_then(|r| r.checksums),
1500 Some(want)
1501 );
1502
1503 Ok(())
1504 }
1505
1506 #[tokio::test]
1507 async fn checksum_mixed_then_precomputed() -> Result {
1508 let mut engine = Checksum {
1509 crc32c: Some(Crc32c::default()),
1510 md5_hash: Some(Md5::default()),
1511 };
1512 engine.update(0, &bytes::Bytes::from_static(VEXING.as_bytes()));
1513 let ck = engine.finalize();
1514
1515 let client = test_builder().build().await?;
1516 let upload = client
1517 .write_object("my-bucket", "my-object", QUICK)
1518 .with_known_md5_hash(ck.md5_hash.clone())
1519 .with_known_crc32c(ck.crc32c.unwrap())
1520 .precompute_checksums()
1521 .await?;
1522 // Note that the checksums do not match the data. This is intentional,
1523 // we are trying to verify that whatever is provided in with_known*()
1524 // is respected.
1525 let want = ck.clone();
1526 assert_eq!(
1527 upload.request.spec.resource.and_then(|r| r.checksums),
1528 Some(want)
1529 );
1530
1531 Ok(())
1532 }
1533
1534 #[tokio::test]
1535 async fn checksum_full_computed_then_md5_precomputed() -> Result {
1536 let mut engine = Checksum {
1537 crc32c: Some(Crc32c::default()),
1538 md5_hash: Some(Md5::default()),
1539 };
1540 engine.update(0, &bytes::Bytes::from_static(VEXING.as_bytes()));
1541 let ck = engine.finalize();
1542
1543 let client = test_builder().build().await?;
1544 let upload = client
1545 .write_object("my-bucket", "my-object", QUICK)
1546 .compute_md5()
1547 .with_known_md5_hash(ck.md5_hash.clone())
1548 .precompute_checksums()
1549 .await?;
1550 // Note that the checksums do not match the data. This is intentional,
1551 // we are trying to verify that whatever is provided in with_known*()
1552 // is respected.
1553 let want = quick_checksum(Checksum {
1554 crc32c: Some(Crc32c::default()),
1555 md5_hash: None,
1556 })
1557 .set_md5_hash(ck.md5_hash.clone());
1558 assert_eq!(
1559 upload.request.spec.resource.and_then(|r| r.checksums),
1560 Some(want)
1561 );
1562
1563 Ok(())
1564 }
1565
1566 #[tokio::test]
1567 async fn checksum_known_crc32_then_computed_md5() -> Result {
1568 let mut engine = Checksum {
1569 crc32c: Some(Crc32c::default()),
1570 md5_hash: Some(Md5::default()),
1571 };
1572 engine.update(0, &bytes::Bytes::from_static(VEXING.as_bytes()));
1573 let ck = engine.finalize();
1574
1575 let client = test_builder().build().await?;
1576 let upload = client
1577 .write_object("my-bucket", "my-object", QUICK)
1578 .with_known_crc32c(ck.crc32c.unwrap())
1579 .compute_md5()
1580 .with_known_md5_hash(ck.md5_hash.clone())
1581 .precompute_checksums()
1582 .await?;
1583 // Note that the checksums do not match the data. This is intentional,
1584 // we are trying to verify that whatever is provided in with_known*()
1585 // is respected.
1586 let want = ck.clone();
1587 assert_eq!(
1588 upload.request.spec.resource.and_then(|r| r.checksums),
1589 Some(want)
1590 );
1591
1592 Ok(())
1593 }
1594
1595 #[tokio::test]
1596 async fn checksum_known_crc32_then_known_md5() -> Result {
1597 let mut engine = Checksum {
1598 crc32c: Some(Crc32c::default()),
1599 md5_hash: Some(Md5::default()),
1600 };
1601 engine.update(0, &bytes::Bytes::from_static(VEXING.as_bytes()));
1602 let ck = engine.finalize();
1603
1604 let client = test_builder().build().await?;
1605 let upload = client
1606 .write_object("my-bucket", "my-object", QUICK)
1607 .with_known_crc32c(ck.crc32c.unwrap())
1608 .with_known_md5_hash(ck.md5_hash.clone())
1609 .precompute_checksums()
1610 .await?;
1611 // Note that the checksums do not match the data. This is intentional,
1612 // we are trying to verify that whatever is provided in with_known*()
1613 // is respected.
1614 let want = ck.clone();
1615 assert_eq!(
1616 upload.request.spec.resource.and_then(|r| r.checksums),
1617 Some(want)
1618 );
1619
1620 Ok(())
1621 }
1622
1623 #[tokio::test]
1624 async fn precompute_checksums_seek_error() -> Result {
1625 let mut source = MockSeekSource::new();
1626 source
1627 .expect_seek()
1628 .once()
1629 .returning(|_| Err(IoError::new(ErrorKind::Deadlock, "test-only")));
1630
1631 let client = test_builder().build().await?;
1632 let err = client
1633 .write_object("my-bucket", "my-object", source)
1634 .precompute_checksums()
1635 .await
1636 .expect_err("seek() returns an error");
1637 assert!(err.is_serialization(), "{err:?}");
1638 assert!(
1639 err.source()
1640 .and_then(|e| e.downcast_ref::<IoError>())
1641 .is_some(),
1642 "{err:?}"
1643 );
1644
1645 Ok(())
1646 }
1647
1648 #[tokio::test]
1649 async fn precompute_checksums_next_error() -> Result {
1650 let mut source = MockSeekSource::new();
1651 source.expect_seek().returning(|_| Ok(()));
1652 let mut seq = mockall::Sequence::new();
1653 source
1654 .expect_next()
1655 .times(3)
1656 .in_sequence(&mut seq)
1657 .returning(|| Some(Ok(bytes::Bytes::new())));
1658 source
1659 .expect_next()
1660 .once()
1661 .in_sequence(&mut seq)
1662 .returning(|| Some(Err(IoError::new(ErrorKind::BrokenPipe, "test-only"))));
1663
1664 let client = test_builder().build().await?;
1665 let err = client
1666 .write_object("my-bucket", "my-object", source)
1667 .precompute_checksums()
1668 .await
1669 .expect_err("seek() returns an error");
1670 assert!(err.is_serialization(), "{err:?}");
1671 assert!(
1672 err.source()
1673 .and_then(|e| e.downcast_ref::<IoError>())
1674 .is_some(),
1675 "{err:?}"
1676 );
1677
1678 Ok(())
1679 }
1680
1681 #[tokio::test]
1682 async fn write_object_with_user_agent() -> Result {
1683 use http::header::USER_AGENT;
1684
1685 let user_agent = "quicker_foxes_lazier_dogs/1.2.3";
1686 let server = Server::run();
1687 server.expect(
1688 Expectation::matching(all_of![
1689 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1690 request::headers(contains((USER_AGENT.as_str(), user_agent))),
1691 request::query(url_decoded(contains(("uploadType", "multipart")))),
1692 ])
1693 .times(1)
1694 .respond_with(status_code(200).body("{}")),
1695 );
1696
1697 let client = Storage::builder()
1698 .with_endpoint(format!("http://{}", server.addr()))
1699 .with_credentials(Anonymous::new().build())
1700 .build()
1701 .await?;
1702 let _ = client
1703 .write_object(
1704 "projects/_/buckets/test-bucket",
1705 "test-object",
1706 "hello world",
1707 )
1708 .with_user_agent(user_agent)
1709 .send_unbuffered()
1710 .await?;
1711
1712 Ok(())
1713 }
1714
1715 #[tokio::test]
1716 async fn write_object_with_quota_project() -> Result {
1717 const PROJECT_NAME: &str = "project_lazy_dog";
1718 let server = Server::run();
1719 let session = server.url("/upload/session/test-only-001");
1720 let path = session.path().to_string();
1721
1722 server.expect(
1723 Expectation::matching(all_of![
1724 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1725 request::headers(contains(("x-goog-user-project", PROJECT_NAME))),
1726 request::query(url_decoded(contains(("uploadType", "resumable")))),
1727 ])
1728 .times(1)
1729 .respond_with(status_code(200).append_header("location", session.to_string())),
1730 );
1731
1732 server.expect(
1733 Expectation::matching(all_of![
1734 request::method_path("PUT", path),
1735 request::headers(contains(("x-goog-user-project", PROJECT_NAME))),
1736 ])
1737 .times(1)
1738 .respond_with(
1739 status_code(200)
1740 .append_header(http::header::CONTENT_TYPE, "application/json")
1741 .body(serde_json::to_string(&crate::model::Object::new()).unwrap()),
1742 ),
1743 );
1744
1745 let client = Storage::builder()
1746 .with_endpoint(format!("http://{}", server.addr()))
1747 .with_credentials(Anonymous::new().build())
1748 .build()
1749 .await?;
1750 let _ = client
1751 .write_object(
1752 "projects/_/buckets/test-bucket",
1753 "test-object",
1754 "hello world",
1755 )
1756 .with_quota_project(PROJECT_NAME)
1757 .with_resumable_upload_threshold(0_usize)
1758 .send_unbuffered()
1759 .await?;
1760
1761 Ok(())
1762 }
1763
1764 #[tokio::test]
1765 async fn debug() -> Result {
1766 let client = test_builder().build().await?;
1767 let upload = client
1768 .write_object("my-bucket", "my-object", "")
1769 .precompute_checksums()
1770 .await;
1771
1772 let fmt = format!("{upload:?}");
1773 ["WriteObject", "inner", "spec", "options", "checksum"]
1774 .into_iter()
1775 .for_each(|text| {
1776 assert!(fmt.contains(text), "expected {text} in {fmt}");
1777 });
1778 Ok(())
1779 }
1780
1781 #[tokio::test]
1782 async fn with_checksum_precomputation_builder() -> Result {
1783 let client = test_builder().build().await?;
1784 let upload = client
1785 .write_object("my-bucket", "my-object", "hello")
1786 .with_checksum_precomputation(true);
1787 assert!(upload.request.checksum_precomputation);
1788
1789 let upload = client
1790 .write_object("my-bucket", "my-object", "hello")
1791 .with_checksum_precomputation(false);
1792 assert!(!upload.request.checksum_precomputation);
1793
1794 let upload = client.write_object("my-bucket", "my-object", "hello");
1795 assert!(upload.request.checksum_precomputation);
1796 Ok(())
1797 }
1798
1799 #[tokio::test]
1800 async fn unbuffered_resumable_precomputes_by_default() -> Result {
1801 let server = Server::run();
1802 let session = server.url("/upload/session/test-only-001");
1803 let path = session.path().to_string();
1804
1805 use base64::{Engine, prelude::BASE64_STANDARD};
1806 let expected_crc = crc32c::crc32c(b"hello world");
1807 let expected_crc_b64 = BASE64_STANDARD.encode(expected_crc.to_be_bytes());
1808
1809 server.expect(
1810 Expectation::matching(all_of![
1811 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1812 request::query(url_decoded(contains(("uploadType", "resumable")))),
1813 request::body(json_decoded(move |body: &serde_json::Value| {
1814 body.get("crc32c").and_then(|v| v.as_str()) == Some(&expected_crc_b64)
1815 })),
1816 ])
1817 .times(1)
1818 .respond_with(status_code(200).append_header("location", session.to_string())),
1819 );
1820
1821 server.expect(
1822 Expectation::matching(all_of![
1823 request::method_path("PUT", path),
1824 request::body("hello world"),
1825 ])
1826 .times(1)
1827 .respond_with(
1828 status_code(200)
1829 .append_header(http::header::CONTENT_TYPE, "application/json")
1830 .body(serde_json::to_string(&crate::model::Object::new()).unwrap()),
1831 ),
1832 );
1833
1834 let client = Storage::builder()
1835 .with_endpoint(format!("http://{}", server.addr()))
1836 .with_credentials(Anonymous::new().build())
1837 .build()
1838 .await?;
1839 let _ = client
1840 .write_object(
1841 "projects/_/buckets/test-bucket",
1842 "test-object",
1843 "hello world",
1844 )
1845 .with_resumable_upload_threshold(0_usize)
1846 .send_unbuffered()
1847 .await?;
1848
1849 Ok(())
1850 }
1851
1852 #[tokio::test]
1853 async fn unbuffered_resumable_with_checksum_precomputation_false() -> Result {
1854 let server = Server::run();
1855 let session = server.url("/upload/session/test-only-002");
1856 let path = session.path().to_string();
1857
1858 server.expect(
1859 Expectation::matching(all_of![
1860 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1861 request::query(url_decoded(contains(("uploadType", "resumable")))),
1862 request::body(json_decoded(|body: &serde_json::Value| {
1863 body.get("crc32c").is_none()
1864 })),
1865 ])
1866 .times(1)
1867 .respond_with(status_code(200).append_header("location", session.to_string())),
1868 );
1869
1870 server.expect(
1871 Expectation::matching(request::method_path("PUT", path))
1872 .times(1)
1873 .respond_with(
1874 status_code(200)
1875 .append_header(http::header::CONTENT_TYPE, "application/json")
1876 .body(serde_json::to_string(&crate::model::Object::new()).unwrap()),
1877 ),
1878 );
1879
1880 let client = Storage::builder()
1881 .with_endpoint(format!("http://{}", server.addr()))
1882 .with_credentials(Anonymous::new().build())
1883 .build()
1884 .await?;
1885 let _ = client
1886 .write_object(
1887 "projects/_/buckets/test-bucket",
1888 "test-object",
1889 "hello world",
1890 )
1891 .with_resumable_upload_threshold(0_usize)
1892 .with_checksum_precomputation(false)
1893 .send_unbuffered()
1894 .await?;
1895
1896 Ok(())
1897 }
1898
1899 #[tokio::test]
1900 async fn unbuffered_single_shot_does_not_precompute_by_default() -> Result {
1901 let server = Server::run();
1902 server.expect(
1903 Expectation::matching(all_of![
1904 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1905 request::query(url_decoded(contains(("uploadType", "multipart")))),
1906 ])
1907 .times(1)
1908 .respond_with(
1909 status_code(200)
1910 .append_header(http::header::CONTENT_TYPE, "application/json")
1911 .body(serde_json::to_string(&crate::model::Object::new()).unwrap()),
1912 ),
1913 );
1914
1915 let client = Storage::builder()
1916 .with_endpoint(format!("http://{}", server.addr()))
1917 .with_credentials(Anonymous::new().build())
1918 .build()
1919 .await?;
1920 let _ = client
1921 .write_object(
1922 "projects/_/buckets/test-bucket",
1923 "test-object",
1924 "hello world",
1925 )
1926 .send_unbuffered()
1927 .await?;
1928
1929 Ok(())
1930 }
1931
1932 #[tokio::test]
1933 async fn unbuffered_single_shot_with_checksum_precomputation_false() -> Result {
1934 let server = Server::run();
1935 server.expect(
1936 Expectation::matching(all_of![
1937 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1938 request::query(url_decoded(contains(("uploadType", "multipart")))),
1939 ])
1940 .times(1)
1941 .respond_with(
1942 status_code(200)
1943 .append_header(http::header::CONTENT_TYPE, "application/json")
1944 .body(serde_json::to_string(&crate::model::Object::new()).unwrap()),
1945 ),
1946 );
1947
1948 let client = Storage::builder()
1949 .with_endpoint(format!("http://{}", server.addr()))
1950 .with_credentials(Anonymous::new().build())
1951 .build()
1952 .await?;
1953 let _ = client
1954 .write_object(
1955 "projects/_/buckets/test-bucket",
1956 "test-object",
1957 "hello world",
1958 )
1959 .with_checksum_precomputation(false)
1960 .send_unbuffered()
1961 .await?;
1962
1963 Ok(())
1964 }
1965
1966 #[tokio::test]
1967 async fn unbuffered_resumable_with_known_crc32c_skips_precompute() -> Result {
1968 let server = Server::run();
1969 let session = server.url("/upload/session/test-only-003");
1970 let path = session.path().to_string();
1971
1972 let known_crc = 123456_u32;
1973 use base64::{Engine, prelude::BASE64_STANDARD};
1974 let expected_crc_b64 = BASE64_STANDARD.encode(known_crc.to_be_bytes());
1975
1976 server.expect(
1977 Expectation::matching(all_of![
1978 request::method_path("POST", "/upload/storage/v1/b/test-bucket/o"),
1979 request::query(url_decoded(contains(("uploadType", "resumable")))),
1980 request::body(json_decoded(move |body: &serde_json::Value| {
1981 body.get("crc32c").and_then(|v| v.as_str()) == Some(&expected_crc_b64)
1982 })),
1983 ])
1984 .times(1)
1985 .respond_with(status_code(200).append_header("location", session.to_string())),
1986 );
1987
1988 server.expect(
1989 Expectation::matching(request::method_path("PUT", path))
1990 .times(1)
1991 .respond_with(
1992 status_code(200)
1993 .append_header(http::header::CONTENT_TYPE, "application/json")
1994 .body(serde_json::to_string(&crate::model::Object::new()).unwrap()),
1995 ),
1996 );
1997
1998 let client = Storage::builder()
1999 .with_endpoint(format!("http://{}", server.addr()))
2000 .with_credentials(Anonymous::new().build())
2001 .build()
2002 .await?;
2003 let _ = client
2004 .write_object(
2005 "projects/_/buckets/test-bucket",
2006 "test-object",
2007 "hello world",
2008 )
2009 .with_known_crc32c(known_crc)
2010 .with_resumable_upload_threshold(0_usize)
2011 .send_unbuffered()
2012 .await?;
2013
2014 Ok(())
2015 }
2016}