Skip to main content

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}