Skip to main content

object_store/
lib.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18#![cfg_attr(docsrs, feature(doc_cfg))]
19#![deny(rustdoc::broken_intra_doc_links, rustdoc::bare_urls, rust_2018_idioms)]
20#![warn(
21    missing_copy_implementations,
22    missing_debug_implementations,
23    missing_docs,
24    clippy::explicit_iter_loop,
25    clippy::future_not_send,
26    clippy::use_self,
27    clippy::clone_on_ref_ptr,
28    unreachable_pub
29)]
30
31//! # object_store
32//!
33//! This crate provides a uniform API for interacting with object
34//! storage services and local files via the [`ObjectStore`]
35//! trait.
36//!
37//! Using this crate, the same binary and code can run in multiple
38//! clouds and local test environments, via a simple runtime
39//! configuration change.
40//!
41//! # Highlights
42//!
43//! 1. A high-performance async API focused on providing a consistent interface
44//!    mirroring that of object stores such as [S3]
45//!
46//! 2. Production quality, leading this crate to be used in large
47//!    scale production systems, such as [crates.io] and [InfluxDB IOx]
48//!
49//! 3. Support for advanced functionality, including atomic, conditional reads
50//!    and writes, vectored IO, bulk deletion, and more...
51//!
52//! 4. Stable and predictable governance via the [Apache Arrow] project
53//!
54//! 5. Small dependency footprint, depending on only a small number of common crates
55//!
56//! Originally developed by [InfluxData] and subsequently donated
57//! to [Apache Arrow].
58//!
59//! [Apache Arrow]: https://arrow.apache.org/
60//! [InfluxData]: https://www.influxdata.com/
61//! [crates.io]: https://github.com/rust-lang/crates.io
62//! [ACID]: https://en.wikipedia.org/wiki/ACID
63//! [S3]: https://aws.amazon.com/s3/
64//!
65//! # APIs
66//!
67//! * [`ObjectStore`]: Core object store API
68//! * [`ObjectStoreExt`]: (*New in 0.13.0*) Extension trait with additional convenience methods
69//!
70//! # Available [`ObjectStore`] Implementations
71//!
72//! By default, this crate provides the following implementations:
73//!
74//! * Memory: [`InMemory`](memory::InMemory)
75//!
76//! Feature flags are used to enable support for other implementations:
77//!
78#![cfg_attr(
79    feature = "fs",
80    doc = "* Local filesystem: [`LocalFileSystem`](local::LocalFileSystem)"
81)]
82#![cfg_attr(
83    feature = "gcp-base",
84    doc = "* [`gcp`]: [Google Cloud Storage](https://cloud.google.com/storage/) support. See [`GoogleCloudStorageBuilder`](gcp::GoogleCloudStorageBuilder)"
85)]
86#![cfg_attr(
87    feature = "aws-base",
88    doc = "* [`aws`]: [Amazon S3](https://aws.amazon.com/s3/). See [`AmazonS3Builder`](aws::AmazonS3Builder)"
89)]
90#![cfg_attr(
91    feature = "azure-base",
92    doc = "* [`azure`]: [Azure Blob Storage](https://azure.microsoft.com/en-gb/services/storage/blobs/). See [`MicrosoftAzureBuilder`](azure::MicrosoftAzureBuilder)"
93)]
94#![cfg_attr(
95    feature = "http-base",
96    doc = "* [`http`]: [HTTP/WebDAV Storage](https://datatracker.ietf.org/doc/html/rfc2518). See [`HttpBuilder`](http::HttpBuilder)"
97)]
98//!
99//! See [Feature Flags](#feature-flags) for the full set of flags.
100//!
101//! # Why not a Filesystem Interface?
102//!
103//! The [`ObjectStore`] interface is designed to mirror the APIs
104//! of object stores and *not* filesystems, and thus has stateless APIs instead
105//! of cursor based interfaces such as [`Read`] or [`Seek`] available in filesystems.
106//!
107//! This design provides the following advantages:
108//!
109//! * All operations are atomic, and readers cannot observe partial and/or failed writes
110//! * Methods map directly to object store APIs, providing both efficiency and predictability
111//! * Abstracts away filesystem and operating system specific quirks, ensuring portability
112//! * Allows for functionality not native to filesystems, such as operation preconditions
113//!   and atomic multipart uploads
114//!
115//! This crate does provide [`BufReader`] and [`BufWriter`] adapters
116//! which provide a more filesystem-like API for working with the
117//! [`ObjectStore`] trait, however, they should be used with care
118//!
119//! [`BufReader`]: buffered::BufReader
120//! [`BufWriter`]: buffered::BufWriter
121//! [`Read`]: std::io::Read
122//! [`Seek`]: std::io::Seek
123//!
124//! # Adapters
125//!
126//! [`ObjectStore`] instances can be composed with various adapters
127//! which add additional functionality:
128//!
129//! * Rate Throttling: [`ThrottleConfig`](throttle::ThrottleConfig)
130//! * Concurrent Request Limit: [`LimitStore`](limit::LimitStore)
131//!
132//! # Configuration System
133//!
134//! This crate provides a configuration system inspired by the APIs exposed by [fsspec],
135//! [PyArrow FileSystem], and [Hadoop FileSystem], allowing creating a [`DynObjectStore`]
136//! from a URL and an optional list of key value pairs. This provides a flexible interface
137//! to support a wide variety of user-defined store configurations, with minimal additional
138//! application complexity.
139//!
140//! ```no_run,ignore-wasm32
141//! # #[cfg(feature = "aws")] {
142//! # use url::Url;
143//! # use object_store::{parse_url, parse_url_opts};
144//! # use object_store::aws::{AmazonS3, AmazonS3Builder};
145//! #
146//! #
147//! // Can manually create a specific store variant using the appropriate builder
148//! let store: AmazonS3 = AmazonS3Builder::from_env()
149//!     .with_bucket_name("my-bucket").build().unwrap();
150//!
151//! // Alternatively can create an ObjectStore from an S3 URL
152//! let url = Url::parse("s3://bucket/path").unwrap();
153//! let (store, path) = parse_url(&url).unwrap();
154//! assert_eq!(path.as_ref(), "path");
155//!
156//! // Potentially with additional options
157//! let (store, path) = parse_url_opts(&url, vec![("aws_access_key_id", "...")]).unwrap();
158//!
159//! // Or with URLs that encode the bucket name in the URL path
160//! let url = Url::parse("https://ACCOUNT_ID.r2.cloudflarestorage.com/bucket/path").unwrap();
161//! let (store, path) = parse_url(&url).unwrap();
162//! assert_eq!(path.as_ref(), "path");
163//! # }
164//! ```
165//!
166//! [PyArrow FileSystem]: https://arrow.apache.org/docs/python/generated/pyarrow.fs.FileSystem.html#pyarrow.fs.FileSystem.from_uri
167//! [fsspec]: https://filesystem-spec.readthedocs.io/en/latest/api.html#fsspec.filesystem
168//! [Hadoop FileSystem]: https://hadoop.apache.org/docs/r3.0.0/api/org/apache/hadoop/fs/FileSystem.html#get-java.net.URI-org.apache.hadoop.conf.Configuration-
169//!
170//! # List objects
171//!
172//! Use the [`ObjectStore::list`] method to iterate over objects in
173//! remote storage or files in the local filesystem:
174//!
175//! ```ignore-wasm32
176//! # use object_store::local::LocalFileSystem;
177//! # use std::sync::Arc;
178//! # use object_store::{path::Path, ObjectStore};
179//! # use futures_util::stream::StreamExt;
180//! # // use LocalFileSystem for example
181//! # fn get_object_store() -> Arc<dyn ObjectStore> {
182//! #   Arc::new(LocalFileSystem::new())
183//! # }
184//! #
185//! # async fn example() {
186//! #
187//! // create an ObjectStore
188//! let object_store: Arc<dyn ObjectStore> = get_object_store();
189//!
190//! // Recursively list all files below the 'data' path.
191//! // 1. On AWS S3 this would be the 'data/' prefix
192//! // 2. On a local filesystem, this would be the 'data' directory
193//! let prefix = Path::from("data");
194//!
195//! // Get an `async` stream of Metadata objects:
196//! let mut list_stream = object_store.list(Some(&prefix));
197//!
198//! // Print a line about each object
199//! while let Some(meta) = list_stream.next().await.transpose().unwrap() {
200//!     println!("Name: {}, size: {}", meta.location, meta.size);
201//! }
202//! # }
203//! ```
204//!
205//! Which will print out something like the following:
206//!
207//! ```text
208//! Name: data/file01.parquet, size: 112832
209//! Name: data/file02.parquet, size: 143119
210//! Name: data/child/file03.parquet, size: 100
211//! ...
212//! ```
213//!
214//! # Fetch objects
215//!
216//! Use the [`ObjectStoreExt::get`] / [`ObjectStore::get_opts`] method to fetch the data bytes
217//! from remote storage or files in the local filesystem as a stream.
218//!
219//! ```ignore-wasm32
220//! # use futures_util::TryStreamExt;
221//! # use object_store::local::LocalFileSystem;
222//! # use std::sync::Arc;
223//! #  use bytes::Bytes;
224//! # use object_store::{path::Path, ObjectStore, ObjectStoreExt, GetResult};
225//! # fn get_object_store() -> Arc<dyn ObjectStore> {
226//! #   Arc::new(LocalFileSystem::new())
227//! # }
228//! #
229//! # async fn example() {
230//! #
231//! // Create an ObjectStore
232//! let object_store: Arc<dyn ObjectStore> = get_object_store();
233//!
234//! // Retrieve a specific file
235//! let path = Path::from("data/file01.parquet");
236//!
237//! // Fetch just the file metadata
238//! let meta = object_store.head(&path).await.unwrap();
239//! println!("{meta:?}");
240//!
241//! // Fetch the object including metadata
242//! let result: GetResult = object_store.get(&path).await.unwrap();
243//! assert_eq!(result.meta, meta);
244//!
245//! // Buffer the entire object in memory
246//! let object: Bytes = result.bytes().await.unwrap();
247//! assert_eq!(object.len() as u64, meta.size);
248//!
249//! // Alternatively stream the bytes from object storage
250//! let stream = object_store.get(&path).await.unwrap().into_stream();
251//!
252//! // Count the '0's using `try_fold` from `TryStreamExt` trait
253//! let num_zeros = stream
254//!     .try_fold(0, |acc, bytes| async move {
255//!         Ok(acc + bytes.iter().filter(|b| **b == 0).count())
256//!     }).await.unwrap();
257//!
258//! println!("Num zeros in {} is {}", path, num_zeros);
259//! # }
260//! ```
261//!
262//! # Put Object
263//!
264//! Use the [`ObjectStoreExt::put`] method to atomically write data.
265//!
266//! To upload large objects without buffering them entirely in [`PutPayload`],
267//! see [Multipart Upload](#multipart-upload) below.
268//!
269//! ```ignore-wasm32
270//! # use object_store::local::LocalFileSystem;
271//! # use object_store::{ObjectStore, ObjectStoreExt, PutPayload};
272//! # use std::sync::Arc;
273//! # use object_store::path::Path;
274//! # fn get_object_store() -> Arc<dyn ObjectStore> {
275//! #   Arc::new(LocalFileSystem::new())
276//! # }
277//! # async fn put() {
278//! #
279//! let object_store: Arc<dyn ObjectStore> = get_object_store();
280//! let path = Path::from("data/file1");
281//! let payload = PutPayload::from_static(b"hello");
282//! object_store.put(&path, payload).await.unwrap();
283//! # }
284//! ```
285//!
286//! # Multipart Upload
287//!
288//! Use the [`ObjectStoreExt::put_multipart`] / [`ObjectStore::put_multipart_opts`] method to atomically write a large
289//! amount of data in multiple parts, without buffering the entire object in memory.
290//! [`WriteMultipart`] uploads fixed size parts in parallel as data is written.
291//!
292//! [`BufWriter`](buffered::BufWriter) provides an [`AsyncWrite`]
293//! interface that automatically picks between a single and multipart upload
294//! based on the amount of data written.
295//!
296//! [`AsyncWrite`]: tokio::io::AsyncWrite
297//! [`BufWriter`]: buffered::BufWriter
298//!
299//! ```ignore-wasm32
300//! # use object_store::local::LocalFileSystem;
301//! # use object_store::{ObjectStore, ObjectStoreExt, WriteMultipart};
302//! # use std::sync::Arc;
303//! # use bytes::Bytes;
304//! # use tokio::io::AsyncWriteExt;
305//! # use object_store::path::Path;
306//! # fn get_object_store() -> Arc<dyn ObjectStore> {
307//! #   Arc::new(LocalFileSystem::new())
308//! # }
309//! # async fn multi_upload() {
310//! #
311//! let object_store: Arc<dyn ObjectStore> = get_object_store();
312//! let path = Path::from("data/large_file");
313//! let upload =  object_store.put_multipart(&path).await.unwrap();
314//! let mut write = WriteMultipart::new(upload);
315//! write.write(b"hello");
316//! write.finish().await.unwrap();
317//! # }
318//! ```
319//!
320//! # Vectored Read
321//!
322//! A common pattern, especially when reading structured datasets, is to need to fetch
323//! multiple, potentially non-contiguous, ranges of a particular object.
324//!
325//! [`ObjectStore::get_ranges`] provides an efficient way to perform such vectored IO, and will
326//! automatically coalesce adjacent ranges into an appropriate number of parallel requests.
327//!
328//! ```ignore-wasm32
329//! # use object_store::local::LocalFileSystem;
330//! # use object_store::ObjectStore;
331//! # use std::sync::Arc;
332//! # use bytes::Bytes;
333//! # use tokio::io::AsyncWriteExt;
334//! # use object_store::path::Path;
335//! # fn get_object_store() -> Arc<dyn ObjectStore> {
336//! #   Arc::new(LocalFileSystem::new())
337//! # }
338//! # async fn multi_upload() {
339//! #
340//! let object_store: Arc<dyn ObjectStore> = get_object_store();
341//! let path = Path::from("data/large_file");
342//! let ranges = object_store.get_ranges(&path, &[90..100, 400..600, 0..10]).await.unwrap();
343//! assert_eq!(ranges.len(), 3);
344//! assert_eq!(ranges[0].len(), 10);
345//! # }
346//! ```
347//!
348//! To retrieve ranges from a versioned object, use [`ObjectStore::get_opts`] by specifying the range in the [`GetOptions`].
349//!
350//! ```ignore-wasm32
351//! # use object_store::local::LocalFileSystem;
352//! # use object_store::ObjectStore;
353//! # use object_store::GetOptions;
354//! # use std::sync::Arc;
355//! # use bytes::Bytes;
356//! # use tokio::io::AsyncWriteExt;
357//! # use object_store::path::Path;
358//! # fn get_object_store() -> Arc<dyn ObjectStore> {
359//! #   Arc::new(LocalFileSystem::new())
360//! # }
361//! # async fn get_range_with_options() {
362//! #
363//! let object_store: Arc<dyn ObjectStore> = get_object_store();
364//! let path = Path::from("data/large_file");
365//! let ranges = vec![90..100, 400..600, 0..10];
366//! for range in ranges {
367//!     let opts = GetOptions::default().with_range(Some(range));
368//!     let data = object_store.get_opts(&path, opts).await.unwrap();
369//!     // Do something with the data
370//! }
371//! # }
372//! ``````
373//!
374//! # Vectored Write
375//!
376//! When writing data it is often the case that the size of the output is not known ahead of time.
377//!
378//! A common approach to handling this is to bump-allocate a `Vec`, whereby the underlying
379//! allocation is repeatedly reallocated, each time doubling the capacity. The performance of
380//! this is suboptimal as reallocating memory will often involve copying it to a new location.
381//!
382//! Fortunately, as [`PutPayload`] does not require memory regions to be contiguous, it is
383//! possible to instead allocate memory in chunks and avoid bump allocating. [`PutPayloadMut`]
384//! encapsulates this approach
385//!
386//! ```ignore-wasm32
387//! # use object_store::local::LocalFileSystem;
388//! # use object_store::{ObjectStore, ObjectStoreExt, PutPayloadMut};
389//! # use std::sync::Arc;
390//! # use bytes::Bytes;
391//! # use tokio::io::AsyncWriteExt;
392//! # use object_store::path::Path;
393//! # fn get_object_store() -> Arc<dyn ObjectStore> {
394//! #   Arc::new(LocalFileSystem::new())
395//! # }
396//! # async fn multi_upload() {
397//! #
398//! let object_store: Arc<dyn ObjectStore> = get_object_store();
399//! let path = Path::from("data/large_file");
400//! let mut buffer = PutPayloadMut::new().with_block_size(8192);
401//! for _ in 0..22 {
402//!     buffer.extend_from_slice(&[0; 1024]);
403//! }
404//! let payload = buffer.freeze();
405//!
406//! // Payload consists of 3 separate 8KB allocations
407//! assert_eq!(payload.as_ref().len(), 3);
408//! assert_eq!(payload.as_ref()[0].len(), 8192);
409//! assert_eq!(payload.as_ref()[1].len(), 8192);
410//! assert_eq!(payload.as_ref()[2].len(), 6144);
411//!
412//! object_store.put(&path, payload).await.unwrap();
413//! # }
414//! ```
415//!
416//! # Conditional Fetch
417//!
418//! More complex object retrieval can be supported by [`ObjectStore::get_opts`].
419//!
420//! For example, efficiently refreshing a cache without re-fetching the entire object
421//! data if the object hasn't been modified.
422//!
423//! ```
424//! # use std::collections::btree_map::Entry;
425//! # use std::collections::HashMap;
426//! # use object_store::{GetOptions, GetResult, ObjectStore, ObjectStoreExt, Result, Error};
427//! # use std::sync::Arc;
428//! # use std::time::{Duration, Instant};
429//! # use bytes::Bytes;
430//! # use tokio::io::AsyncWriteExt;
431//! # use object_store::path::Path;
432//! struct CacheEntry {
433//!     /// Data returned by last request
434//!     data: Bytes,
435//!     /// ETag identifying the object returned by the server
436//!     e_tag: String,
437//!     /// Instant of last refresh
438//!     refreshed_at: Instant,
439//! }
440//!
441//! /// Example cache that checks entries after 10 seconds for a new version
442//! struct Cache {
443//!     entries: HashMap<Path, CacheEntry>,
444//!     store: Arc<dyn ObjectStore>,
445//! }
446//!
447//! impl Cache {
448//!     pub async fn get(&mut self, path: &Path) -> Result<Bytes> {
449//!         Ok(match self.entries.get_mut(path) {
450//!             Some(e) => match e.refreshed_at.elapsed() < Duration::from_secs(10) {
451//!                 true => e.data.clone(), // Return cached data
452//!                 false => { // Check if remote version has changed
453//!                     let opts = GetOptions::new().with_if_none_match(Some(e.e_tag.clone()));
454//!                     match self.store.get_opts(&path, opts).await {
455//!                         Ok(d) => e.data = d.bytes().await?,
456//!                         Err(Error::NotModified { .. }) => {} // Data has not changed
457//!                         Err(e) => return Err(e),
458//!                     };
459//!                     e.refreshed_at = Instant::now();
460//!                     e.data.clone()
461//!                 }
462//!             },
463//!             None => { // Not cached, fetch data
464//!                 let get = self.store.get(&path).await?;
465//!                 let e_tag = get.meta.e_tag.clone();
466//!                 let data = get.bytes().await?;
467//!                 if let Some(e_tag) = e_tag {
468//!                     let entry = CacheEntry {
469//!                         e_tag,
470//!                         data: data.clone(),
471//!                         refreshed_at: Instant::now(),
472//!                     };
473//!                     self.entries.insert(path.clone(), entry);
474//!                 }
475//!                 data
476//!             }
477//!         })
478//!     }
479//! }
480//! ```
481//!
482//! # Conditional Put
483//!
484//! The default behaviour when writing data is to upsert any existing object at the given path,
485//! overwriting any previous value. More complex behaviours can be achieved using [`PutMode`], and
486//! can be used to build [Optimistic Concurrency Control] based transactions. This facilitates
487//! building metadata catalogs, such as [Apache Iceberg] or [Delta Lake], directly on top of object
488//! storage, without relying on a separate DBMS.
489//!
490//! ```
491//! # use object_store::{Error, ObjectStore, ObjectStoreExt, PutMode, UpdateVersion};
492//! # use std::sync::Arc;
493//! # use bytes::Bytes;
494//! # use tokio::io::AsyncWriteExt;
495//! # use object_store::memory::InMemory;
496//! # use object_store::path::Path;
497//! # fn get_object_store() -> Arc<dyn ObjectStore> {
498//! #   Arc::new(InMemory::new())
499//! # }
500//! # fn do_update(b: Bytes) -> Bytes {b}
501//! # async fn conditional_put() {
502//! let store = get_object_store();
503//! let path = Path::from("test");
504//!
505//! // Perform a conditional update on path
506//! loop {
507//!     // Perform get request
508//!     let r = store.get(&path).await.unwrap();
509//!
510//!     // Save version information fetched
511//!     let version = UpdateVersion {
512//!         e_tag: r.meta.e_tag.clone(),
513//!         version: r.meta.version.clone(),
514//!     };
515//!
516//!     // Compute new version of object contents
517//!     let new = do_update(r.bytes().await.unwrap());
518//!
519//!     // Attempt to commit transaction
520//!     match store.put_opts(&path, new.into(), PutMode::Update(version).into()).await {
521//!         Ok(_) => break, // Successfully committed
522//!         Err(Error::Precondition { .. }) => continue, // Object has changed, try again
523//!         Err(e) => panic!("{e}")
524//!     }
525//! }
526//! # }
527//! ```
528//!
529//! [Optimistic Concurrency Control]: https://en.wikipedia.org/wiki/Optimistic_concurrency_control
530//! [Apache Iceberg]: https://iceberg.apache.org/
531//! [Delta Lake]: https://delta.io/
532//!
533//! # Feature Flags
534//!
535//! The feature set is layered so that you can pick an object store
536//! implementation, its HTTP transport, and its cryptography provider
537//! independently:
538//!
539//! * `cloud-base` shared cloud implementation (XML/JSON parsing,
540//!   credentials, retry, etc.) and intentionally does *not* depend on
541//!   `reqwest` or a cryptography provider.
542//! * `reqwest` enables the built-in [`reqwest`]-based [`HttpConnector`].
543//! * `aws-lc-rs` and `ring` each provide a bundled [`client::CryptoProvider`].
544//! * `<provider>-base` (`aws-base`, `azure-base`, `gcp-base`, `http-base`)
545//!   adds the implementation specific logic on top of `cloud-base` without pulling in
546//!   `reqwest` or a cryptography provider.
547//! * `<provider>` (`aws`, `azure`, `gcp`, `http`) is the batteries-included
548//!   feature for `<provider>-base` + `reqwest` (with `rustls`) + the default
549//!   `aws-lc-rs` cryptography provider, and is the typical choice.
550//!
551//! ## Implementation specific features
552//!
553//! | Feature | Enables | Notes |
554//! | --- | --- | --- |
555//! | `aws` | `aws-base` + `reqwest` + `aws-lc-rs` | Amazon S3 with the built-in HTTP transport. |
556//! | `azure` | `azure-base` + `reqwest` + `aws-lc-rs` | Azure Blob Storage with the built-in HTTP transport. |
557//! | `gcp` | `gcp-base` + `reqwest` + `aws-lc-rs` | Google Cloud Storage with the built-in HTTP transport. |
558//! | `http` | `http-base` + `reqwest` + `aws-lc-rs` | HTTP/WebDAV with the built-in HTTP transport. |
559//! | `aws-base` |  | S3 without `reqwest` or crypto; supply your own [`HttpConnector`] and [`client::CryptoProvider`]. |
560//! | `azure-base` |  | Azure without `reqwest` or crypto; supply your own [`HttpConnector`] and [`client::CryptoProvider`]. |
561//! | `gcp-base` |  | GCS without `reqwest` or crypto; supply your own [`HttpConnector`] and [`client::CryptoProvider`]. |
562//! | `http-base` |  | HTTP/WebDAV without `reqwest`; supply your own [`HttpConnector`]. |
563//!
564//! ## Transport and crypto features
565//!
566//! | Feature | Description |
567//! | --- | --- |
568//! | `reqwest` | Enables the [`reqwest`]-based [`HttpConnector`]. Enabled automatically by `aws`, `azure`, `gcp`, and `http`. |
569//! | `aws-lc-rs` | Bundled [`aws-lc-rs`]-based [`client::CryptoProvider`]. The default for the batteries-included provider features. |
570//! | `ring` | Bundled [`ring`]-based [`client::CryptoProvider`], e.g. for WASM targets. |
571//! | `cloud-base` | Shared cloud-provider implementation. Enabled automatically by `*-base` features; usually not enabled directly. |
572//!
573//! ## Other features
574//!
575//! | Feature | Description |
576//! | --- | --- |
577//! | `fs` *(default)* | Local filesystem store via [`LocalFileSystem`](local::LocalFileSystem). |
578//! | `tokio` | Enables Tokio-based utilities such as [`BufReader`](buffered::BufReader) and [`BufWriter`](buffered::BufWriter). Pulled in automatically by `fs` and the `*-base` features. |
579//! | `integration` | Exposes the [`integration`] module, a reusable test suite for verifying custom [`ObjectStore`] implementations. Not API-stable. |
580//!
581//! ## Selecting a `reqwest` TLS backend
582//!
583//! `reqwest` needs a TLS backend to compile, so whenever you enable the `reqwest` feature directly
584//! you must also enable one of `reqwest`'s TLS features:
585//!
586//! | reqwest feature | TLS stack | Notes |
587//! | --- | --- | --- |
588//! | `reqwest/rustls` | [rustls] with [`aws-lc-rs`] | enables `aws-lc-rs`. This is what `aws`/`azure`/`gcp`/`http` enable. |
589//! | `reqwest/native-tls` | the platform's native TLS (OpenSSL / SChannel / Secure Transport) | enables neither `rustls` nor `aws-lc-rs`. |
590//! | `reqwest/rustls-no-provider` | [rustls] with no bundled provider | enables neither provider; you must install one at runtime, e.g. `rustls::crypto::ring::default_provider().install_default()`. |
591//!
592//! ## Feature examples
593//!
594//! S3 implementation only; user provides the HTTP connector and crypto provider:
595//! ```toml
596//! object_store = { default-features = false, features = ["aws-base"] }
597//! ```
598//!
599//! S3 implementation + `reqwest` + `aws-lc-rs` signing (equivalent to the `aws` feature):
600//! ```toml
601//! object_store = { default-features = false, features = ["aws-base", "reqwest", "reqwest/rustls", "aws-lc-rs"] }
602//! ```
603//!
604//! S3 implementation + `reqwest` with native TLS + `ring` signing (no `aws-lc-rs` in the dependency tree):
605//! ```toml
606//! object_store = { default-features = false, features = ["aws-base", "reqwest", "reqwest/native-tls", "ring"] }
607//! ```
608//!
609//! [rustls]: https://crates.io/crates/rustls/
610//!
611//! # Cryptography
612//!
613//! Request signing (e.g. AWS SigV4 or GCP service-account signing) requires a
614//! [`client::CryptoProvider`]. The `aws`, `gcp`, and `azure` features
615//! use [`aws-lc-rs`], matching `reqwest`'s default so that applications do not end up with
616//! two crypto stacks.
617//!
618//! If you wish to use [`ring`] (e.g. to support WASM targets), use the
619//! `*-base` feature flags, e.g. `aws-base`, and then enable the `ring` feature.
620//!
621//! When enabling the `aws-lc-rs` feature without the built-in `reqwest`
622//! transport (e.g. `aws-base` + `aws-lc-rs` with a custom [`HttpConnector`]),
623//! you must also select an [`aws-lc-rs`] backend, as `object_store` does not
624//! pick one for you.
625//! For example enable `aws-lc-rs/aws-lc-sys` (or `fips` / `non-fips`).
626//!
627//! If both `ring` and `aws-lc-rs` are enabled, `aws-lc-rs` is used by default.
628//!
629//! You can also implement a custom [`client::CryptoProvider`] to use your own cryptographic library.
630//!
631//! This signing provider is independent of the TLS crypto provider used by the
632//! built-in `reqwest` transport — see
633//! [Selecting a `reqwest` TLS backend](#selecting-a-reqwest-tls-backend). The
634//! only combination that needs the provider registered manually (e.g.
635//! `rustls::crypto::ring::default_provider().install_default()` in your `main`)
636//! is `reqwest/rustls-no-provider`; `reqwest/rustls` and `reqwest/native-tls`
637//! configure their TLS stack automatically.
638//!
639//! [`aws-lc-rs`]: https://crates.io/crates/aws-lc-rs/
640//! [`ring`]: https://crates.io/crates/ring/
641//!
642//! # TLS Certificates
643//!
644//! Stores that use HTTPS/TLS (this is true for most cloud stores) can choose how certificates are validated.
645//!
646//! By default [`rustls-platform-verifier`] is used to verify certificates using the system's certificate
647//! facilities. Alternatively, this functionality can be disabled using
648//! [`ClientOptions::with_no_system_certificates`] and certificates manually registered using
649//! [`ClientOptions::with_root_certificate`].
650//!
651//! These could be a custom CA chain, or alternatively an alternative trust store, e.g. [`webpki-roots`].
652//!
653//! ```ignore-wasm32
654//! # #[cfg(feature = "aws")] {
655//! use object_store::{ClientOptions, Certificate};
656//!
657//! let mut options = ClientOptions::default().with_no_system_certificates(true);
658//! for root_cert in webpki_root_certs::TLS_SERVER_ROOT_CERTS {
659//!     options = options.with_root_certificate(Certificate::from_der(root_cert.as_ref()).unwrap());
660//! }
661//! # }
662//! ```
663//!
664//! [CA]: https://en.wikipedia.org/wiki/Certificate_authority
665//! [`rustls-platform-verifier`]: https://crates.io/crates/rustls-platform-verifier/
666//! [`webpki-roots`]: https://crates.io/crates/webpki-roots
667//!
668//! # Customizing HTTP Clients
669//!
670//! Many [`ObjectStore`] implementations permit customization of the HTTP client via
671//! the [`HttpConnector`] trait and utilities in the [`client`] module.
672//! Examples include injecting custom HTTP headers or using an alternate
673//! tokio Runtime for I/O requests. To replace `reqwest` entirely (rather than
674//! tweak the bundled transport) see [Disabling `reqwest`](#disabling-reqwest).
675//!
676//! [`HttpConnector`]: client::HttpConnector
677//!
678//! # Disabling `reqwest`
679//!
680//! The `aws`, `azure`, `gcp`, and `http` features each bundle a
681//! [`reqwest`]-based HTTP transport, which is the right choice for most
682//! applications. If you would rather supply your own HTTP client — for example
683//! to share an existing client, to target a platform where `reqwest` does not
684//! compile (such as `wasm32-wasip1`), or to keep `reqwest` out of your
685//! dependency tree — use the matching `*-base` feature and provide an
686//! [`HttpConnector`](client::HttpConnector) at builder time.
687//!
688//! Remember to disable the default features so that `fs` (and its transitive
689//! dependencies) is not pulled in:
690//!
691//! ```toml
692//! [dependencies]
693//! object_store = { version = "0.13", default-features = false, features = ["aws-base"] }
694//! ```
695//!
696//! ```ignore
697//! use object_store::aws::AmazonS3Builder;
698//!
699//! let store = AmazonS3Builder::from_env()
700//!     // `my_connector` is your own `impl HttpConnector`
701//!     .with_http_connector(my_connector)
702//!     .build()?;
703//! ```
704//!
705//! See [Feature Flags](#feature-flags) above for the full set of flags.
706//!
707//! [`reqwest`]: https://crates.io/crates/reqwest
708
709#[cfg(feature = "aws-base")]
710pub mod aws;
711#[cfg(feature = "azure-base")]
712pub mod azure;
713#[cfg(feature = "tokio")]
714pub mod buffered;
715#[cfg(not(target_arch = "wasm32"))]
716pub mod chunked;
717pub mod delimited;
718#[cfg(feature = "gcp-base")]
719pub mod gcp;
720#[cfg(feature = "http-base")]
721pub mod http;
722#[cfg(feature = "tokio")]
723pub mod limit;
724#[cfg(all(feature = "fs", not(target_arch = "wasm32")))]
725pub mod local;
726pub mod memory;
727pub mod path;
728pub mod prefix;
729pub mod registry;
730#[cfg(any(feature = "aws-base", feature = "azure-base", feature = "gcp-base"))]
731pub mod retry;
732#[cfg(feature = "cloud-base")]
733pub mod signer;
734#[cfg(feature = "tokio")]
735pub mod throttle;
736
737#[cfg(feature = "cloud-base")]
738pub mod client;
739
740#[cfg(feature = "cloud-base")]
741pub use client::{
742    ClientConfigKey, ClientOptions, CredentialProvider, StaticCredentialProvider,
743    backoff::BackoffConfig, retry::RetryConfig,
744};
745
746#[cfg(all(
747    feature = "cloud-base",
748    feature = "reqwest",
749    not(target_arch = "wasm32")
750))]
751pub use client::Certificate;
752
753#[cfg(feature = "cloud-base")]
754mod config;
755
756mod tags;
757
758pub use tags::TagSet;
759
760pub mod list;
761pub mod multipart;
762mod parse;
763mod payload;
764mod upload;
765mod util;
766
767mod attributes;
768
769#[cfg(any(feature = "integration", test))]
770pub mod integration;
771
772pub use attributes::*;
773
774pub use parse::{ObjectStoreScheme, parse_url, parse_url_opts};
775pub use payload::*;
776pub use upload::*;
777pub use util::{GetRange, OBJECT_STORE_COALESCE_DEFAULT, coalesce_ranges, collect_bytes};
778
779// Re-export HTTP types used in public API
780pub use ::http::{Extensions, HeaderMap, HeaderValue};
781
782use crate::path::Path;
783#[cfg(all(feature = "fs", not(target_arch = "wasm32")))]
784use crate::util::maybe_spawn_blocking;
785use async_trait::async_trait;
786use bytes::Bytes;
787use chrono::{DateTime, Utc};
788use futures_util::{StreamExt, TryStreamExt, stream::BoxStream};
789use std::fmt::{Debug, Formatter};
790use std::ops::Range;
791use std::sync::Arc;
792
793/// An alias for a dynamically dispatched object store implementation.
794pub type DynObjectStore = dyn ObjectStore;
795
796/// Id type for multipart uploads.
797pub type MultipartId = String;
798
799/// Universal API for object store services.
800///
801/// See the [module-level documentation](crate) for a high level overview and
802/// examples. See [`ObjectStoreExt`] for additional convenience methods.
803///
804/// # Contract
805/// This trait is a contract between object store _implementations_
806/// (e.g. providers, wrappers) and the `object_store` crate itself. It is
807/// intended to be the minimum API required for an object store.
808///
809/// The [`ObjectStoreExt`] acts as an API/contract between `object_store` and
810/// the store _users_ and provides additional methods that may be simpler to use
811/// but overlap in functionality with [`ObjectStore`].
812///
813/// # Clone
814/// If a store implements [`Clone`], that will only clone the handle to the underlying data. It will NOT clone/fork the
815/// actual key-value data. Hence, the cloned instance and the original instance share the same state.
816///
817/// # Minimal Default Implementations
818/// There are only a few default implementations for methods in this trait by
819/// design. This was different from versions prior to `0.13.0`, which had many
820/// more default implementations. Default implementations are convenient for
821/// users, but error-prone for implementors as they require keeping the
822/// convenience APIs correctly in sync.
823///
824/// As of version 0.13.0, most methods on [`ObjectStore`] must be implemented, and
825/// the convenience methods have been moved to the [`ObjectStoreExt`] trait as
826/// described above. See [#385] for more details.
827///
828/// [#385]: https://github.com/apache/arrow-rs-object-store/issues/385
829///
830/// # Wrappers
831/// If you wrap an [`ObjectStore`] -- e.g. to add observability -- you SHOULD
832/// implement all trait methods. This ensures that defaults implementations
833/// that are overwritten by the wrapped store are also used by the wrapper.
834/// For example:
835///
836/// ```ignore
837/// struct MyStore {
838///     ...
839/// }
840///
841/// #[async_trait]
842/// impl ObjectStore for MyStore {
843///     // implement custom ranges handling
844///     async fn get_ranges(
845///         &self,
846///         location: &Path,
847///         ranges: &[Range<u64>],
848///     ) -> Result<Vec<Bytes>> {
849///         ...
850///     }
851///
852///     ...
853/// }
854///
855/// struct Wrapper {
856///     inner: Arc<dyn ObjectStore>,
857/// }
858///
859/// #[async_trait]
860/// #[deny(clippy::missing_trait_methods)]
861/// impl ObjectStore for Wrapper {
862///     // If we would not implement this method,
863///     // we would get the trait default and not
864///     // use the actual implementation of `inner`.
865///     async fn get_ranges(
866///         &self,
867///         location: &Path,
868///         ranges: &[Range<u64>],
869///     ) -> Result<Vec<Bytes>> {
870///         ...
871///     }
872///
873///     ...
874/// }
875/// ```
876///
877/// To automatically detect this issue, use
878/// [`#[deny(clippy::missing_trait_methods)]`](https://rust-lang.github.io/rust-clippy/master/index.html#missing_trait_methods).
879///
880/// # Upgrade Guide for 0.13.0
881///
882/// Upgrading to object_store 0.13.0 from an earlier version typically involves:
883///
884/// 1. Add a `use` for [`ObjectStoreExt`] to solve the error
885///
886/// ```text
887/// error[E0599]: no method named `put` found for reference `&dyn object_store::ObjectStore` in the current scope
888///    --> datafusion/datasource/src/url.rs:993:14
889/// ```
890///
891/// 2. Remove any (now) redundant implementations (such as `ObjectStore::put`) from any
892///   `ObjectStore` implementations to resolve the error
893///
894/// ```text
895/// error[E0407]: method `put` is not a member of trait `ObjectStore`
896///     --> datafusion/datasource/src/url.rs:1103:9
897///      |
898/// ```
899///
900/// 3. Convert `ObjectStore::delete` to [`ObjectStore::delete_stream`] (see documentation
901///    on that method for details and examples)
902///
903/// 4. Combine `ObjectStore::copy` and `ObjectStore::copy_if_not_exists` implementations into
904///    [`ObjectStore::copy_opts`] (see documentation on that method for details and examples)
905///
906/// 5. Update `object_store::Error::NotImplemented` to include the name of the missing method
907///
908/// For example, change instances of
909/// ```text
910/// object_store::Error::NotImplemented
911/// ```
912/// to
913/// ```
914/// object_store::Error::NotImplemented {
915///    operation: "put".to_string(),
916///    implementer: "RequestCountingObjectStore".to_string(),
917///  };
918/// ```
919///
920#[async_trait]
921pub trait ObjectStore: std::fmt::Display + Send + Sync + Debug + 'static {
922    /// Save the provided `payload` to `location` with the given options
923    ///
924    /// The operation is guaranteed to be atomic, it will either successfully
925    /// write the entirety of `payload` to `location`, or fail. No clients
926    /// should be able to observe a partially written object
927    ///
928    /// To upload large objects without buffering them entirely in memory, use
929    /// the multipart API: [`ObjectStore::put_multipart_opts`]
930    async fn put_opts(
931        &self,
932        location: &Path,
933        payload: PutPayload,
934        opts: PutOptions,
935    ) -> Result<PutResult>;
936
937    /// Perform a multipart upload with options
938    ///
939    /// Client should prefer [`ObjectStore::put_opts`] for small payloads, as streaming uploads
940    /// typically require multiple separate requests. See [`MultipartUpload`] for more information
941    ///
942    /// See also [`BufWriter`](buffered::BufWriter) for an interface that
943    /// automatically picks between a single and multipart upload based on the
944    /// amount of data written.
945    ///
946    /// For more advanced multipart uploads see [`MultipartStore`](multipart::MultipartStore)
947    ///
948    /// [`BufWriter`]: buffered::BufWriter
949    async fn put_multipart_opts(
950        &self,
951        location: &Path,
952        opts: PutMultipartOptions,
953    ) -> Result<Box<dyn MultipartUpload>>;
954
955    /// Perform a get request with options
956    ///
957    /// ## Example
958    ///
959    /// This example uses a basic local filesystem object store to get an object with a specific etag.
960    /// On the local filesystem, supplying an invalid etag will error.
961    /// Versioned object stores will return the specified object version, if it exists.
962    ///
963    /// ```ignore-wasm32
964    /// # use object_store::local::LocalFileSystem;
965    /// # use tempfile::tempdir;
966    /// # use object_store::{path::Path, ObjectStore, ObjectStoreExt, GetOptions};
967    /// async fn get_opts_example() {
968    ///     let tmp = tempdir().unwrap();
969    ///     let store = LocalFileSystem::new_with_prefix(tmp.path()).unwrap();
970    ///     let location = Path::from("example.txt");
971    ///     let content = b"Hello, Object Store!";
972    ///
973    ///     // Put the object into the store
974    ///     store
975    ///         .put(&location, content.as_ref().into())
976    ///         .await
977    ///         .expect("Failed to put object");
978    ///
979    ///     // Get the object from the store to figure out the right etag
980    ///     let result: object_store::GetResult = store.get(&location).await.expect("Failed to get object");
981    ///
982    ///     let etag = result.meta.e_tag.expect("ETag should be present");
983    ///
984    ///     // Get the object from the store with range and etag
985    ///     let bytes = store
986    ///         .get_opts(
987    ///             &location,
988    ///             GetOptions::new()
989    ///                 .with_if_match(Some(etag.clone())),
990    ///         )
991    ///         .await
992    ///         .expect("Failed to get object with range and etag")
993    ///         .bytes()
994    ///         .await
995    ///         .expect("Failed to read bytes");
996    ///
997    ///     println!(
998    ///         "Retrieved with ETag {}: {}",
999    ///         etag,
1000    ///         String::from_utf8_lossy(&bytes)
1001    ///     );
1002    ///
1003    ///     // Show that if the etag does not match, we get an error
1004    ///     let wrong_etag = "wrong-etag".to_string();
1005    ///     match store
1006    ///         .get_opts(
1007    ///             &location,
1008    ///             GetOptions::new().with_if_match(Some(wrong_etag))
1009    ///         )
1010    ///         .await
1011    ///     {
1012    ///         Ok(_) => println!("Unexpectedly succeeded with wrong ETag"),
1013    ///         Err(e) => println!("On a non-versioned object store, getting an invalid ETag ('wrong-etag') results in an error as expected: {}", e),
1014    ///     }
1015    /// }
1016    /// ```
1017    ///
1018    /// To retrieve a range of bytes from a versioned object, specify the range in the [`GetOptions`] supplied to this method.
1019    ///
1020    /// ```ignore-wasm32
1021    /// # use object_store::local::LocalFileSystem;
1022    /// # use tempfile::tempdir;
1023    /// # use object_store::{path::Path, ObjectStore, ObjectStoreExt, GetOptions};
1024    /// async fn get_opts_range_example() {
1025    ///     let tmp = tempdir().unwrap();
1026    ///     let store = LocalFileSystem::new_with_prefix(tmp.path()).unwrap();
1027    ///     let location = Path::from("example.txt");
1028    ///     let content = b"Hello, Object Store!";
1029    ///
1030    ///     // Put the object into the store
1031    ///     store
1032    ///         .put(&location, content.as_ref().into())
1033    ///         .await
1034    ///         .expect("Failed to put object");
1035    ///
1036    ///     // Get the object from the store to figure out the right etag
1037    ///     let result: object_store::GetResult = store.get(&location).await.expect("Failed to get object");
1038    ///
1039    ///     let etag = result.meta.e_tag.expect("ETag should be present");
1040    ///
1041    ///     // Get the object from the store with range and etag
1042    ///     let bytes = store
1043    ///         .get_opts(
1044    ///             &location,
1045    ///             GetOptions::new()
1046    ///                 .with_range(Some(0..5))
1047    ///                 .with_if_match(Some(etag.clone())),
1048    ///         )
1049    ///         .await
1050    ///         .expect("Failed to get object with range and etag")
1051    ///         .bytes()
1052    ///         .await
1053    ///         .expect("Failed to read bytes");
1054    ///
1055    ///     println!(
1056    ///         "Retrieved range [0-5] with ETag {}: {}",
1057    ///         etag,
1058    ///         String::from_utf8_lossy(&bytes)
1059    ///     );
1060    ///
1061    ///     // Show that if the etag does not match, we get an error
1062    ///     let wrong_etag = "wrong-etag".to_string();
1063    ///     match store
1064    ///         .get_opts(
1065    ///             &location,
1066    ///             GetOptions::new().with_range(Some(0..5)).with_if_match(Some(wrong_etag))
1067    ///         )
1068    ///         .await
1069    ///     {
1070    ///         Ok(_) => println!("Unexpectedly succeeded with wrong ETag"),
1071    ///         Err(e) => println!("On a non-versioned object store, getting an invalid ETag ('wrong-etag') results in an error as expected: {}", e),
1072    ///     }
1073    /// }
1074    /// ```
1075    async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult>;
1076
1077    /// Return the bytes that are stored at the specified location
1078    /// in the given byte ranges
1079    async fn get_ranges(&self, location: &Path, ranges: &[Range<u64>]) -> Result<Vec<Bytes>> {
1080        coalesce_ranges(
1081            ranges,
1082            |range| self.get_range(location, range),
1083            OBJECT_STORE_COALESCE_DEFAULT,
1084        )
1085        .await
1086    }
1087
1088    /// Delete all the objects at the specified locations
1089    ///
1090    /// When supported, this method will use bulk operations that delete more
1091    /// than one object per a request. Otherwise, the implementation may call
1092    /// the single object delete method for each location.
1093    ///
1094    /// # Bulk Delete Support
1095    ///
1096    /// The following backends support native bulk delete operations:
1097    ///
1098    /// - **AWS (S3)**: Uses the native [DeleteObjects] API with batches of up to 1000 objects
1099    /// - **Azure**: Uses the native [Blob Batch] API with batches of up to 256 objects
1100    ///
1101    /// The following backends use concurrent individual delete operations:
1102    ///
1103    /// - **GCP**: Performs individual delete requests with up to 10 concurrent operations
1104    /// - **HTTP**: Performs individual delete requests with up to 10 concurrent operations
1105    /// - **Local**: Performs individual file deletions with up to 10 concurrent operations
1106    /// - **Memory**: Performs individual in-memory deletions sequentially
1107    ///
1108    /// [DeleteObjects]: https://docs.aws.amazon.com/AmazonS3/latest/API/API_DeleteObjects.html
1109    /// [Blob Batch]: https://learn.microsoft.com/en-us/rest/api/storageservices/blob-batch
1110    ///
1111    /// The returned stream yields the results of the delete operations in the
1112    /// same order as the input locations. However, some errors will be from
1113    /// an overall call to a bulk delete operation, and not from a specific
1114    /// location.
1115    ///
1116    /// If the object did not exist, the result may be an error or a success,
1117    /// depending on the behavior of the underlying store. For example, local
1118    /// filesystems, GCP, and Azure return an error, while S3 and in-memory will
1119    /// return Ok. If it is an error, it will be [`Error::NotFound`].
1120    ///
1121    /// ```ignore-wasm32
1122    /// # use futures_util::{StreamExt, TryStreamExt};
1123    /// # use object_store::local::LocalFileSystem;
1124    /// # async fn example() -> Result<(), Box<dyn std::error::Error>> {
1125    /// # let root = tempfile::TempDir::new().unwrap();
1126    /// # let store = LocalFileSystem::new_with_prefix(root.path()).unwrap();
1127    /// # use object_store::{ObjectStore, ObjectStoreExt, ObjectMeta};
1128    /// # use object_store::path::Path;
1129    /// # use futures_util::{StreamExt, TryStreamExt};
1130    /// #
1131    /// // Create two objects
1132    /// store.put(&Path::from("foo"), "foo".into()).await?;
1133    /// store.put(&Path::from("bar"), "bar".into()).await?;
1134    ///
1135    /// // List object
1136    /// let locations = store.list(None).map_ok(|m| m.location).boxed();
1137    ///
1138    /// // Delete them
1139    /// store.delete_stream(locations).try_collect::<Vec<Path>>().await?;
1140    /// # Ok(())
1141    /// # }
1142    /// # let rt = tokio::runtime::Builder::new_current_thread().build().unwrap();
1143    /// # rt.block_on(example()).unwrap();
1144    /// ```
1145    ///
1146    /// Note: Before version 0.13, `delete_stream` has a default implementation
1147    /// that deletes each object with up to 10 concurrent requests. This default
1148    /// behavior has been removed, and each implementation must now provide its
1149    /// own `delete_stream` implementation explicitly. The following example
1150    /// shows how to implement `delete_stream` to get the previous default
1151    /// behavior.
1152    ///
1153    /// ```
1154    /// # use async_trait::async_trait;
1155    /// # use futures_util::stream::{BoxStream, StreamExt};
1156    /// # use object_store::path::Path;
1157    /// # use object_store::{
1158    /// #     CopyOptions, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore,
1159    /// #     PutMultipartOptions, PutOptions, PutPayload, PutResult, Result,
1160    /// # };
1161    /// # use std::fmt;
1162    /// # use std::fmt::Debug;
1163    /// # use std::sync::Arc;
1164    /// #
1165    /// # struct ExampleClient;
1166    /// #
1167    /// # impl ExampleClient {
1168    /// #     async fn delete(&self, _path: &Path) -> Result<()> {
1169    /// #         Ok(())
1170    /// #     }
1171    /// # }
1172    /// #
1173    /// # struct ExampleStore {
1174    /// #     client: Arc<ExampleClient>,
1175    /// # }
1176    /// #
1177    /// # impl Debug for ExampleStore {
1178    /// #     fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1179    /// #         write!(f, "ExampleStore")
1180    /// #     }
1181    /// # }
1182    /// #
1183    /// # impl fmt::Display for ExampleStore {
1184    /// #     fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1185    /// #         write!(f, "ExampleStore")
1186    /// #     }
1187    /// # }
1188    /// #
1189    /// # #[async_trait]
1190    /// # impl ObjectStore for ExampleStore {
1191    /// #     async fn put_opts(&self, _: &Path, _: PutPayload, _: PutOptions) -> Result<PutResult> {
1192    /// #         todo!()
1193    /// #     }
1194    /// #
1195    /// #     async fn put_multipart_opts(
1196    /// #         &self,
1197    /// #         _: &Path,
1198    /// #         _: PutMultipartOptions,
1199    /// #     ) -> Result<Box<dyn MultipartUpload>> {
1200    /// #         todo!()
1201    /// #     }
1202    /// #
1203    /// #     async fn get_opts(&self, _: &Path, _: GetOptions) -> Result<GetResult> {
1204    /// #         todo!()
1205    /// #     }
1206    /// #
1207    /// fn delete_stream(
1208    ///     &self,
1209    ///     locations: BoxStream<'static, Result<Path>>,
1210    /// ) -> BoxStream<'static, Result<Path>> {
1211    ///     let client = Arc::clone(&self.client);
1212    ///     locations
1213    ///         .map(move |location| {
1214    ///             let client = Arc::clone(&client);
1215    ///             async move {
1216    ///                 let location = location?;
1217    ///                 client.delete(&location).await?;
1218    ///                 Ok(location)
1219    ///             }
1220    ///         })
1221    ///         .buffered(10)
1222    ///         .boxed()
1223    /// }
1224    /// #
1225    /// #     fn list(&self, _: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
1226    /// #         todo!()
1227    /// #     }
1228    /// #
1229    /// #     async fn list_with_delimiter(&self, _: Option<&Path>) -> Result<ListResult> {
1230    /// #         todo!()
1231    /// #     }
1232    /// #
1233    /// #     async fn copy_opts(&self, _: &Path, _: &Path, _: CopyOptions) -> Result<()> {
1234    /// #         todo!()
1235    /// #     }
1236    /// # }
1237    /// #
1238    /// # async fn example() {
1239    /// #     let store = ExampleStore { client: Arc::new(ExampleClient) };
1240    /// #     let paths = futures_util::stream::iter(vec![Ok(Path::from("foo")), Ok(Path::from("bar"))]).boxed();
1241    /// #     let results = store.delete_stream(paths).collect::<Vec<_>>().await;
1242    /// #     assert_eq!(results.len(), 2);
1243    /// #     assert_eq!(results[0].as_ref().unwrap(), &Path::from("foo"));
1244    /// #     assert_eq!(results[1].as_ref().unwrap(), &Path::from("bar"));
1245    /// # }
1246    /// #
1247    /// # let rt = tokio::runtime::Builder::new_current_thread().build().unwrap();
1248    /// # rt.block_on(example());
1249    /// ```
1250    fn delete_stream(
1251        &self,
1252        locations: BoxStream<'static, Result<Path>>,
1253    ) -> BoxStream<'static, Result<Path>>;
1254
1255    /// List all the objects with the given prefix.
1256    ///
1257    /// Prefixes are evaluated on a path segment basis, i.e. `foo/bar` is a prefix of `foo/bar/x` but not of
1258    /// `foo/bar_baz/x`. List is recursive, i.e. `foo/bar/more/x` will be included.
1259    ///
1260    /// Note: the order of returned [`ObjectMeta`] is not guaranteed
1261    ///
1262    /// For more advanced listing see [`PaginatedListStore`](list::PaginatedListStore)
1263    fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>>;
1264
1265    /// List all the objects with the given prefix and a location greater than `offset`
1266    ///
1267    /// Some stores, such as S3 and GCS, may be able to push `offset` down to reduce
1268    /// the number of network requests required.
1269    ///
1270    /// This returns an exclusive offset, i.e. objects at exactly `offset` will not be included.
1271    ///
1272    /// Note: the order of returned [`ObjectMeta`] is not guaranteed
1273    ///
1274    /// For more advanced listing see [`PaginatedListStore`](list::PaginatedListStore)
1275    fn list_with_offset(
1276        &self,
1277        prefix: Option<&Path>,
1278        offset: &Path,
1279    ) -> BoxStream<'static, Result<ObjectMeta>> {
1280        let offset = offset.clone();
1281        self.list(prefix)
1282            .try_filter(move |f| futures_util::future::ready(f.location > offset))
1283            .boxed()
1284    }
1285
1286    /// List objects with the given prefix and an implementation specific
1287    /// delimiter. Returns common prefixes (directories) in addition to object
1288    /// metadata.
1289    ///
1290    /// Prefixes are evaluated on a path segment basis, i.e. `foo/bar` is a prefix of `foo/bar/x` but not of
1291    /// `foo/bar_baz/x`. List is not recursive, i.e. `foo/bar/more/x` will not be included.
1292    async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult>;
1293
1294    /// Copy an object from one path to another in the same object store.
1295    async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()>;
1296
1297    /// Move an object from one path to another in the same object store.
1298    ///
1299    /// By default, this is implemented as a copy and then delete source. It may not
1300    /// check when deleting source that it was the same object that was originally copied.
1301    async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
1302        let RenameOptions {
1303            target_mode,
1304            extensions,
1305        } = options;
1306        let copy_mode = match target_mode {
1307            RenameTargetMode::Overwrite => CopyMode::Overwrite,
1308            RenameTargetMode::Create => CopyMode::Create,
1309        };
1310        let copy_options = CopyOptions {
1311            mode: copy_mode,
1312            extensions,
1313        };
1314        self.copy_opts(from, to, copy_options).await?;
1315        self.delete(from).await?;
1316        Ok(())
1317    }
1318}
1319
1320macro_rules! as_ref_impl {
1321    ($type:ty) => {
1322        #[async_trait]
1323        #[deny(clippy::missing_trait_methods)]
1324        impl<T: ObjectStore + ?Sized> ObjectStore for $type {
1325            async fn put_opts(
1326                &self,
1327                location: &Path,
1328                payload: PutPayload,
1329                opts: PutOptions,
1330            ) -> Result<PutResult> {
1331                self.as_ref().put_opts(location, payload, opts).await
1332            }
1333
1334            async fn put_multipart_opts(
1335                &self,
1336                location: &Path,
1337                opts: PutMultipartOptions,
1338            ) -> Result<Box<dyn MultipartUpload>> {
1339                self.as_ref().put_multipart_opts(location, opts).await
1340            }
1341
1342            async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
1343                self.as_ref().get_opts(location, options).await
1344            }
1345
1346            async fn get_ranges(
1347                &self,
1348                location: &Path,
1349                ranges: &[Range<u64>],
1350            ) -> Result<Vec<Bytes>> {
1351                self.as_ref().get_ranges(location, ranges).await
1352            }
1353
1354            fn delete_stream(
1355                &self,
1356                locations: BoxStream<'static, Result<Path>>,
1357            ) -> BoxStream<'static, Result<Path>> {
1358                self.as_ref().delete_stream(locations)
1359            }
1360
1361            fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
1362                self.as_ref().list(prefix)
1363            }
1364
1365            fn list_with_offset(
1366                &self,
1367                prefix: Option<&Path>,
1368                offset: &Path,
1369            ) -> BoxStream<'static, Result<ObjectMeta>> {
1370                self.as_ref().list_with_offset(prefix, offset)
1371            }
1372
1373            async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
1374                self.as_ref().list_with_delimiter(prefix).await
1375            }
1376
1377            async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
1378                self.as_ref().copy_opts(from, to, options).await
1379            }
1380
1381            async fn rename_opts(
1382                &self,
1383                from: &Path,
1384                to: &Path,
1385                options: RenameOptions,
1386            ) -> Result<()> {
1387                self.as_ref().rename_opts(from, to, options).await
1388            }
1389        }
1390    };
1391}
1392
1393as_ref_impl!(Arc<T>);
1394as_ref_impl!(Box<T>);
1395
1396/// Extension trait for [`ObjectStore`] with convenience functions.
1397///
1398/// See the [module-level documentation](crate) for a high level overview and
1399/// examples. See "contract" section within the [`ObjectStore`] documentation
1400/// for more reasoning.
1401///
1402/// # Implementation
1403/// You MUST NOT implement this trait yourself. It is automatically implemented for all [`ObjectStore`] implementations.
1404pub trait ObjectStoreExt: ObjectStore {
1405    /// Save the provided bytes to the specified location
1406    ///
1407    /// The operation is guaranteed to be atomic, it will either successfully
1408    /// write the entirety of `payload` to `location`, or fail. No clients
1409    /// should be able to observe a partially written object
1410    ///
1411    /// Note the entire `payload` is buffered in memory. To upload large objects
1412    /// without buffering them entirely in memory, use the multipart API:
1413    /// [`ObjectStoreExt::put_multipart`], [`WriteMultipart`], or
1414    /// [`BufWriter`](buffered::BufWriter)
1415    fn put(&self, location: &Path, payload: PutPayload) -> impl Future<Output = Result<PutResult>>;
1416
1417    /// Perform a multipart upload
1418    ///
1419    /// Client should prefer [`ObjectStoreExt::put`] for small payloads, as streaming uploads
1420    /// typically require multiple separate requests. See [`MultipartUpload`] for more information
1421    ///
1422    /// For more advanced multipart uploads see [`MultipartStore`](multipart::MultipartStore)
1423    fn put_multipart(
1424        &self,
1425        location: &Path,
1426    ) -> impl Future<Output = Result<Box<dyn MultipartUpload>>>;
1427
1428    /// Return the bytes that are stored at the specified location.
1429    ///
1430    /// ## Example
1431    ///
1432    /// This example uses a basic local filesystem object store to get an object.
1433    ///
1434    /// ```ignore-wasm32
1435    /// # use object_store::local::LocalFileSystem;
1436    /// # use tempfile::tempdir;
1437    /// # use object_store::{path::Path, ObjectStore, ObjectStoreExt};
1438    /// async fn get_example() {
1439    ///     let tmp = tempdir().unwrap();
1440    ///     let store = LocalFileSystem::new_with_prefix(tmp.path()).unwrap();
1441    ///     let location = Path::from("example.txt");
1442    ///     let content = b"Hello, Object Store!";
1443    ///
1444    ///     // Put the object into the store
1445    ///     store
1446    ///         .put(&location, content.as_ref().into())
1447    ///         .await
1448    ///         .expect("Failed to put object");
1449    ///
1450    ///     // Get the object from the store
1451    ///     let get_result = store.get(&location).await.expect("Failed to get object");
1452    ///     let bytes = get_result.bytes().await.expect("Failed to read bytes");
1453    ///     println!("Retrieved content: {}", String::from_utf8_lossy(&bytes));
1454    /// }
1455    /// ```
1456    fn get(&self, location: &Path) -> impl Future<Output = Result<GetResult>>;
1457
1458    /// Return the bytes that are stored at the specified location
1459    /// in the given byte range.
1460    ///
1461    /// See [`GetRange::Bounded`] for more details on how `range` gets interpreted.
1462    ///
1463    /// To retrieve a range of bytes from a versioned object, use [`ObjectStore::get_opts`] by specifying the range in the [`GetOptions`].
1464    ///
1465    /// ## Examples
1466    ///
1467    /// This example uses a basic local filesystem object store to get a byte range from an object.
1468    ///
1469    /// ```ignore-wasm32
1470    /// # use object_store::local::LocalFileSystem;
1471    /// # use tempfile::tempdir;
1472    /// # use object_store::{path::Path, ObjectStore, ObjectStoreExt};
1473    /// async fn get_range_example() {
1474    ///     let tmp = tempdir().unwrap();
1475    ///     let store = LocalFileSystem::new_with_prefix(tmp.path()).unwrap();
1476    ///     let location = Path::from("example.txt");
1477    ///     let content = b"Hello, Object Store!";
1478    ///
1479    ///     // Put the object into the store
1480    ///     store
1481    ///         .put(&location, content.as_ref().into())
1482    ///         .await
1483    ///         .expect("Failed to put object");
1484    ///
1485    ///     // Get the object from the store
1486    ///     let bytes = store
1487    ///         .get_range(&location, 0..5)
1488    ///         .await
1489    ///         .expect("Failed to get object");
1490    ///     println!("Retrieved range [0-5]: {}", String::from_utf8_lossy(&bytes));
1491    /// }
1492    /// ```
1493    fn get_range(&self, location: &Path, range: Range<u64>) -> impl Future<Output = Result<Bytes>>;
1494
1495    /// Return the metadata for the specified location
1496    fn head(&self, location: &Path) -> impl Future<Output = Result<ObjectMeta>>;
1497
1498    /// Delete the object at the specified location.
1499    fn delete(&self, location: &Path) -> impl Future<Output = Result<()>>;
1500
1501    /// Copy an object from one path to another in the same object store.
1502    ///
1503    /// If there exists an object at the destination, it will be overwritten.
1504    fn copy(&self, from: &Path, to: &Path) -> impl Future<Output = Result<()>>;
1505
1506    /// Copy an object from one path to another, only if destination is empty.
1507    ///
1508    /// Will return an error if the destination already has an object.
1509    ///
1510    /// Performs an atomic operation if the underlying object storage supports it.
1511    /// If atomic operations are not supported by the underlying object storage (like S3)
1512    /// it will return an error.
1513    fn copy_if_not_exists(&self, from: &Path, to: &Path) -> impl Future<Output = Result<()>>;
1514
1515    /// Move an object from one path to another in the same object store.
1516    ///
1517    /// By default, this is implemented as a copy and then delete source. It may not
1518    /// check when deleting source that it was the same object that was originally copied.
1519    ///
1520    /// If there exists an object at the destination, it will be overwritten.
1521    fn rename(&self, from: &Path, to: &Path) -> impl Future<Output = Result<()>>;
1522
1523    /// Move an object from one path to another in the same object store.
1524    ///
1525    /// Will return an error if the destination already has an object.
1526    fn rename_if_not_exists(&self, from: &Path, to: &Path) -> impl Future<Output = Result<()>>;
1527}
1528
1529impl<T> ObjectStoreExt for T
1530where
1531    T: ObjectStore + ?Sized,
1532{
1533    async fn put(&self, location: &Path, payload: PutPayload) -> Result<PutResult> {
1534        self.put_opts(location, payload, PutOptions::default())
1535            .await
1536    }
1537
1538    async fn put_multipart(&self, location: &Path) -> Result<Box<dyn MultipartUpload>> {
1539        self.put_multipart_opts(location, PutMultipartOptions::default())
1540            .await
1541    }
1542
1543    async fn get(&self, location: &Path) -> Result<GetResult> {
1544        self.get_opts(location, GetOptions::default()).await
1545    }
1546
1547    async fn get_range(&self, location: &Path, range: Range<u64>) -> Result<Bytes> {
1548        let options = GetOptions::new().with_range(Some(range));
1549        self.get_opts(location, options).await?.bytes().await
1550    }
1551
1552    async fn head(&self, location: &Path) -> Result<ObjectMeta> {
1553        let options = GetOptions::new().with_head(true);
1554        Ok(self.get_opts(location, options).await?.meta)
1555    }
1556
1557    async fn delete(&self, location: &Path) -> Result<()> {
1558        let location = location.clone();
1559        let mut stream =
1560            self.delete_stream(futures_util::stream::once(async move { Ok(location) }).boxed());
1561        let _path = stream.try_next().await?.ok_or_else(|| Error::Generic {
1562            store: "ext",
1563            source: "`delete_stream` with one location should yield once but didn't".into(),
1564        })?;
1565        if stream.next().await.is_some() {
1566            Err(Error::Generic {
1567                store: "ext",
1568                source: "`delete_stream` with one location expected to yield exactly once, but yielded more than once".into(),
1569            })
1570        } else {
1571            Ok(())
1572        }
1573    }
1574
1575    async fn copy(&self, from: &Path, to: &Path) -> Result<()> {
1576        let options = CopyOptions::new().with_mode(CopyMode::Overwrite);
1577        self.copy_opts(from, to, options).await
1578    }
1579
1580    async fn copy_if_not_exists(&self, from: &Path, to: &Path) -> Result<()> {
1581        let options = CopyOptions::new().with_mode(CopyMode::Create);
1582        self.copy_opts(from, to, options).await
1583    }
1584
1585    async fn rename(&self, from: &Path, to: &Path) -> Result<()> {
1586        let options = RenameOptions::new().with_target_mode(RenameTargetMode::Overwrite);
1587        self.rename_opts(from, to, options).await
1588    }
1589
1590    async fn rename_if_not_exists(&self, from: &Path, to: &Path) -> Result<()> {
1591        let options = RenameOptions::new().with_target_mode(RenameTargetMode::Create);
1592        self.rename_opts(from, to, options).await
1593    }
1594}
1595
1596/// Result of a list call that includes objects, prefixes (directories) and a
1597/// token for the next set of results. Individual result sets may be limited to
1598/// 1,000 objects based on the underlying object storage's limitations.
1599#[derive(Debug)]
1600pub struct ListResult {
1601    /// Prefixes that are common (like directories)
1602    pub common_prefixes: Vec<Path>,
1603    /// Object metadata for the listing
1604    pub objects: Vec<ObjectMeta>,
1605    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
1606    /// that need to return context-specific information (like cache status) from trait methods.
1607    ///
1608    /// HTTP-backed stores in this crate populate this with the extensions of the HTTP
1609    /// response, allowing custom HTTP middleware to propagate information to callers.
1610    /// Where a result is assembled from multiple paginated requests, the extensions of
1611    /// each response are merged, with those of later responses taking precedence.
1612    pub extensions: Extensions,
1613}
1614
1615/// The metadata that describes an object.
1616#[derive(Debug, Clone, PartialEq, Eq)]
1617pub struct ObjectMeta {
1618    /// The full path to the object
1619    pub location: Path,
1620    /// The last modified time
1621    pub last_modified: DateTime<Utc>,
1622    /// The size in bytes of the object.
1623    ///
1624    /// Note this is not `usize` as `object_store` supports 32-bit architectures such as WASM
1625    pub size: u64,
1626    /// The unique identifier for the object
1627    ///
1628    /// <https://datatracker.ietf.org/doc/html/rfc9110#name-etag>
1629    pub e_tag: Option<String>,
1630    /// A version indicator for this object
1631    pub version: Option<String>,
1632}
1633
1634/// Options for a get request, such as range
1635#[derive(Debug, Default, Clone)]
1636pub struct GetOptions {
1637    /// Request will succeed if the `ObjectMeta::e_tag` matches
1638    /// otherwise returning [`Error::Precondition`]
1639    ///
1640    /// See <https://datatracker.ietf.org/doc/html/rfc9110#name-if-match>
1641    ///
1642    /// Examples:
1643    ///
1644    /// ```text
1645    /// If-Match: "xyzzy"
1646    /// If-Match: "xyzzy", "r2d2xxxx", "c3piozzzz"
1647    /// If-Match: *
1648    /// ```
1649    pub if_match: Option<String>,
1650    /// Request will succeed if the `ObjectMeta::e_tag` does not match
1651    /// otherwise returning [`Error::NotModified`]
1652    ///
1653    /// See <https://datatracker.ietf.org/doc/html/rfc9110#section-13.1.2>
1654    ///
1655    /// Examples:
1656    ///
1657    /// ```text
1658    /// If-None-Match: "xyzzy"
1659    /// If-None-Match: "xyzzy", "r2d2xxxx", "c3piozzzz"
1660    /// If-None-Match: *
1661    /// ```
1662    pub if_none_match: Option<String>,
1663    /// Request will succeed if the object has been modified since
1664    ///
1665    /// <https://datatracker.ietf.org/doc/html/rfc9110#section-13.1.3>
1666    pub if_modified_since: Option<DateTime<Utc>>,
1667    /// Request will succeed if the object has not been modified since
1668    /// otherwise returning [`Error::Precondition`]
1669    ///
1670    /// Some stores, such as S3, will only return `NotModified` for exact
1671    /// timestamp matches, instead of for any timestamp greater than or equal.
1672    ///
1673    /// <https://datatracker.ietf.org/doc/html/rfc9110#section-13.1.4>
1674    pub if_unmodified_since: Option<DateTime<Utc>>,
1675    /// Request transfer of only the specified range of bytes
1676    /// otherwise returning [`Error::NotModified`]
1677    ///
1678    /// <https://datatracker.ietf.org/doc/html/rfc9110#name-range>
1679    pub range: Option<GetRange>,
1680    /// Request a particular object version
1681    pub version: Option<String>,
1682    /// Request transfer of no content
1683    ///
1684    /// <https://datatracker.ietf.org/doc/html/rfc9110#name-head>
1685    pub head: bool,
1686    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
1687    /// that need to pass context-specific information (like tracing spans) via trait methods.
1688    ///
1689    /// These extensions are ignored entirely by backends offered through this crate.
1690    pub extensions: Extensions,
1691}
1692
1693impl GetOptions {
1694    /// Returns an error if the modification conditions on this request are not satisfied
1695    ///
1696    /// <https://datatracker.ietf.org/doc/html/rfc7232#section-6>
1697    pub fn check_preconditions(&self, meta: &ObjectMeta) -> Result<()> {
1698        // The use of the invalid etag "*" means no ETag is equivalent to never matching
1699        let etag = meta.e_tag.as_deref().unwrap_or("*");
1700        let last_modified = meta.last_modified;
1701
1702        if let Some(m) = &self.if_match {
1703            if m != "*" && m.split(',').map(str::trim).all(|x| x != etag) {
1704                return Err(Error::Precondition {
1705                    path: meta.location.to_string(),
1706                    source: format!("{etag} does not match {m}").into(),
1707                });
1708            }
1709        } else if let Some(date) = self.if_unmodified_since {
1710            if last_modified > date {
1711                return Err(Error::Precondition {
1712                    path: meta.location.to_string(),
1713                    source: format!("{date} < {last_modified}").into(),
1714                });
1715            }
1716        }
1717
1718        if let Some(m) = &self.if_none_match {
1719            if m == "*" || m.split(',').map(str::trim).any(|x| x == etag) {
1720                return Err(Error::NotModified {
1721                    path: meta.location.to_string(),
1722                    source: format!("{etag} matches {m}").into(),
1723                });
1724            }
1725        } else if let Some(date) = self.if_modified_since {
1726            if last_modified <= date {
1727                return Err(Error::NotModified {
1728                    path: meta.location.to_string(),
1729                    source: format!("{date} >= {last_modified}").into(),
1730                });
1731            }
1732        }
1733        Ok(())
1734    }
1735
1736    /// Create a new [`GetOptions`]
1737    pub fn new() -> Self {
1738        Self::default()
1739    }
1740
1741    /// Sets the `if_match` condition.
1742    ///
1743    /// See [`GetOptions::if_match`]
1744    #[must_use]
1745    pub fn with_if_match(mut self, etag: Option<impl Into<String>>) -> Self {
1746        self.if_match = etag.map(Into::into);
1747        self
1748    }
1749
1750    /// Sets the `if_none_match` condition.
1751    ///
1752    /// See [`GetOptions::if_none_match`]
1753    #[must_use]
1754    pub fn with_if_none_match(mut self, etag: Option<impl Into<String>>) -> Self {
1755        self.if_none_match = etag.map(Into::into);
1756        self
1757    }
1758
1759    /// Sets the `if_modified_since` condition.
1760    ///
1761    /// See [`GetOptions::if_modified_since`]
1762    #[must_use]
1763    pub fn with_if_modified_since(mut self, dt: Option<impl Into<DateTime<Utc>>>) -> Self {
1764        self.if_modified_since = dt.map(Into::into);
1765        self
1766    }
1767
1768    /// Sets the `if_unmodified_since` condition.
1769    ///
1770    /// See [`GetOptions::if_unmodified_since`]
1771    #[must_use]
1772    pub fn with_if_unmodified_since(mut self, dt: Option<impl Into<DateTime<Utc>>>) -> Self {
1773        self.if_unmodified_since = dt.map(Into::into);
1774        self
1775    }
1776
1777    /// Sets the `range` condition.
1778    ///
1779    /// See [`GetOptions::range`]
1780    #[must_use]
1781    pub fn with_range(mut self, range: Option<impl Into<GetRange>>) -> Self {
1782        self.range = range.map(Into::into);
1783        self
1784    }
1785
1786    /// Sets the `version` condition.
1787    ///
1788    /// See [`GetOptions::version`]
1789    #[must_use]
1790    pub fn with_version(mut self, version: Option<impl Into<String>>) -> Self {
1791        self.version = version.map(Into::into);
1792        self
1793    }
1794
1795    /// Sets the `head` condition.
1796    ///
1797    /// See [`GetOptions::head`]
1798    #[must_use]
1799    pub fn with_head(mut self, head: impl Into<bool>) -> Self {
1800        self.head = head.into();
1801        self
1802    }
1803
1804    /// Sets the `extensions` condition.
1805    ///
1806    /// See [`GetOptions::extensions`]
1807    #[must_use]
1808    pub fn with_extensions(mut self, extensions: Extensions) -> Self {
1809        self.extensions = extensions;
1810        self
1811    }
1812}
1813
1814/// Result for a get request
1815#[derive(Debug)]
1816pub struct GetResult {
1817    /// The [`GetResultPayload`]
1818    pub payload: GetResultPayload,
1819    /// The [`ObjectMeta`] for this object
1820    pub meta: ObjectMeta,
1821    /// The range of bytes returned by this request
1822    ///
1823    /// Note this is not `usize` as `object_store` supports 32-bit architectures such as WASM
1824    pub range: Range<u64>,
1825    /// Additional object attributes
1826    pub attributes: Attributes,
1827    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
1828    /// that need to return context-specific information (like cache status) from trait methods.
1829    ///
1830    /// HTTP-backed stores in this crate populate this with the extensions of the HTTP
1831    /// response, allowing custom HTTP middleware to propagate information to callers.
1832    pub extensions: Extensions,
1833}
1834
1835/// The kind of a [`GetResult`]
1836///
1837/// This special cases the case of a local file, as some systems may
1838/// be able to optimise the case of a file already present on local disk
1839pub enum GetResultPayload {
1840    /// The file, path
1841    #[cfg(all(feature = "fs", not(target_arch = "wasm32")))]
1842    File(std::fs::File, std::path::PathBuf),
1843    /// An opaque stream of bytes
1844    Stream(BoxStream<'static, Result<Bytes>>),
1845}
1846
1847impl Debug for GetResultPayload {
1848    fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
1849        match self {
1850            #[cfg(all(feature = "fs", not(target_arch = "wasm32")))]
1851            Self::File(_, _) => write!(f, "GetResultPayload(File)"),
1852            Self::Stream(_) => write!(f, "GetResultPayload(Stream)"),
1853        }
1854    }
1855}
1856
1857impl GetResult {
1858    /// Collects the data into a [`Bytes`]
1859    pub async fn bytes(self) -> Result<Bytes> {
1860        let len = self.range.end - self.range.start;
1861        match self.payload {
1862            #[cfg(all(feature = "fs", not(target_arch = "wasm32")))]
1863            GetResultPayload::File(mut file, path) => {
1864                maybe_spawn_blocking(move || {
1865                    use crate::local::read_range;
1866
1867                    let buffer = read_range(&mut file, &path, self.range)?;
1868
1869                    Ok(buffer)
1870                })
1871                .await
1872            }
1873            GetResultPayload::Stream(s) => collect_bytes(s, Some(len)).await,
1874        }
1875    }
1876
1877    /// Converts this into a byte stream
1878    ///
1879    /// If the `self.kind` is [`GetResultPayload::File`] will perform chunked reads of the file,
1880    /// otherwise will return the [`GetResultPayload::Stream`].
1881    ///
1882    /// # Tokio Compatibility
1883    ///
1884    /// Tokio discourages performing blocking IO on a tokio worker thread, however,
1885    /// no major operating systems have stable async file APIs. Therefore if called from
1886    /// a tokio context, this will use [`tokio::runtime::Handle::spawn_blocking`] to dispatch
1887    /// IO to a blocking thread pool, much like `tokio::fs` does under-the-hood.
1888    ///
1889    /// If not called from a tokio context, this will perform IO on the current thread with
1890    /// no additional complexity or overheads
1891    pub fn into_stream(self) -> BoxStream<'static, Result<Bytes>> {
1892        match self.payload {
1893            #[cfg(all(feature = "fs", not(target_arch = "wasm32")))]
1894            GetResultPayload::File(file, path) => {
1895                const CHUNK_SIZE: usize = 8 * 1024;
1896                local::chunked_stream(file, path, self.range, CHUNK_SIZE)
1897            }
1898            GetResultPayload::Stream(s) => s,
1899        }
1900    }
1901}
1902
1903/// Configure preconditions for the put operation
1904#[derive(Debug, Clone, PartialEq, Eq, Default)]
1905pub enum PutMode {
1906    /// Perform an atomic write operation, overwriting any object present at the provided path
1907    #[default]
1908    Overwrite,
1909    /// Perform an atomic write operation, returning [`Error::AlreadyExists`] if an
1910    /// object already exists at the provided path
1911    Create,
1912    /// Perform an atomic write operation if the current version of the object matches the
1913    /// provided [`UpdateVersion`], returning [`Error::Precondition`] otherwise
1914    Update(UpdateVersion),
1915}
1916
1917/// Uniquely identifies a version of an object to update
1918///
1919/// Stores will use differing combinations of `e_tag` and `version` to provide conditional
1920/// updates, and it is therefore recommended applications preserve both
1921#[derive(Debug, Clone, PartialEq, Eq)]
1922pub struct UpdateVersion {
1923    /// The unique identifier for the newly created object
1924    ///
1925    /// <https://datatracker.ietf.org/doc/html/rfc9110#name-etag>
1926    pub e_tag: Option<String>,
1927    /// A version indicator for the newly created object
1928    pub version: Option<String>,
1929}
1930
1931impl From<PutResult> for UpdateVersion {
1932    fn from(value: PutResult) -> Self {
1933        Self {
1934            e_tag: value.e_tag,
1935            version: value.version,
1936        }
1937    }
1938}
1939
1940/// Options for a put request
1941#[derive(Debug, Clone, Default)]
1942pub struct PutOptions {
1943    /// Configure the [`PutMode`] for this operation
1944    pub mode: PutMode,
1945    /// Provide a [`TagSet`] for this object
1946    ///
1947    /// Implementations that don't support object tagging should ignore this
1948    pub tags: TagSet,
1949    /// Provide a set of [`Attributes`]
1950    ///
1951    /// Implementations that don't support an attribute should return an error
1952    pub attributes: Attributes,
1953    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
1954    /// that need to pass context-specific information (like tracing spans) via trait methods.
1955    ///
1956    /// These extensions are ignored entirely by backends offered through this crate.
1957    ///
1958    /// They are also excluded from [`PartialEq`] and [`Eq`].
1959    pub extensions: Extensions,
1960}
1961
1962impl PartialEq<Self> for PutOptions {
1963    fn eq(&self, other: &Self) -> bool {
1964        let Self {
1965            mode,
1966            tags,
1967            attributes,
1968            extensions: _,
1969        } = self;
1970        let Self {
1971            mode: other_mode,
1972            tags: other_tags,
1973            attributes: other_attributes,
1974            extensions: _,
1975        } = other;
1976        (mode == other_mode) && (tags == other_tags) && (attributes == other_attributes)
1977    }
1978}
1979
1980impl Eq for PutOptions {}
1981
1982impl From<PutMode> for PutOptions {
1983    fn from(mode: PutMode) -> Self {
1984        Self {
1985            mode,
1986            ..Default::default()
1987        }
1988    }
1989}
1990
1991impl From<TagSet> for PutOptions {
1992    fn from(tags: TagSet) -> Self {
1993        Self {
1994            tags,
1995            ..Default::default()
1996        }
1997    }
1998}
1999
2000impl From<Attributes> for PutOptions {
2001    fn from(attributes: Attributes) -> Self {
2002        Self {
2003            attributes,
2004            ..Default::default()
2005        }
2006    }
2007}
2008
2009// See <https://github.com/apache/arrow-rs-object-store/issues/339>.
2010#[doc(hidden)]
2011#[deprecated(note = "Use PutMultipartOptions", since = "0.12.3")]
2012pub type PutMultipartOpts = PutMultipartOptions;
2013
2014/// Options for [`ObjectStore::put_multipart_opts`]
2015#[derive(Debug, Clone, Default)]
2016pub struct PutMultipartOptions {
2017    /// Provide a [`TagSet`] for this object
2018    ///
2019    /// Implementations that don't support object tagging should ignore this
2020    pub tags: TagSet,
2021    /// Provide a set of [`Attributes`]
2022    ///
2023    /// Implementations that don't support an attribute should return an error
2024    pub attributes: Attributes,
2025    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
2026    /// that need to pass context-specific information (like tracing spans) via trait methods.
2027    ///
2028    /// Cloud backends offered through this crate use extensions installed by methods such as
2029    /// `PutMultipartOptions::with_retry_policy`. Other extensions are ignored.
2030    ///
2031    /// They are also excluded from [`PartialEq`] and [`Eq`].
2032    pub extensions: Extensions,
2033}
2034
2035impl PartialEq<Self> for PutMultipartOptions {
2036    fn eq(&self, other: &Self) -> bool {
2037        let Self {
2038            tags,
2039            attributes,
2040            extensions: _,
2041        } = self;
2042        let Self {
2043            tags: other_tags,
2044            attributes: other_attributes,
2045            extensions: _,
2046        } = other;
2047        (tags == other_tags) && (attributes == other_attributes)
2048    }
2049}
2050
2051impl Eq for PutMultipartOptions {}
2052
2053impl From<TagSet> for PutMultipartOptions {
2054    fn from(tags: TagSet) -> Self {
2055        Self {
2056            tags,
2057            ..Default::default()
2058        }
2059    }
2060}
2061
2062impl From<Attributes> for PutMultipartOptions {
2063    fn from(attributes: Attributes) -> Self {
2064        Self {
2065            attributes,
2066            ..Default::default()
2067        }
2068    }
2069}
2070
2071/// Result for a put request
2072#[derive(Debug, Clone)]
2073pub struct PutResult {
2074    /// The unique identifier for the newly created object
2075    ///
2076    /// <https://datatracker.ietf.org/doc/html/rfc9110#name-etag>
2077    pub e_tag: Option<String>,
2078    /// A version indicator for the newly created object
2079    pub version: Option<String>,
2080    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
2081    /// that need to return context-specific information (like cache status) from trait methods.
2082    ///
2083    /// HTTP-backed stores in this crate populate this with the extensions of the HTTP
2084    /// response, allowing custom HTTP middleware to propagate information to callers.
2085    ///
2086    /// These extensions are excluded from [`PartialEq`] and [`Eq`].
2087    pub extensions: Extensions,
2088}
2089
2090impl PartialEq<Self> for PutResult {
2091    fn eq(&self, other: &Self) -> bool {
2092        let Self {
2093            e_tag,
2094            version,
2095            extensions: _,
2096        } = self;
2097        let Self {
2098            e_tag: other_e_tag,
2099            version: other_version,
2100            extensions: _,
2101        } = other;
2102        (e_tag == other_e_tag) && (version == other_version)
2103    }
2104}
2105
2106impl Eq for PutResult {}
2107
2108/// Configure preconditions for the copy operation
2109#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
2110pub enum CopyMode {
2111    /// Perform an atomic write operation, overwriting any object present at the provided path
2112    #[default]
2113    Overwrite,
2114    /// Perform an atomic write operation, returning [`Error::AlreadyExists`] if an
2115    /// object already exists at the provided path
2116    Create,
2117}
2118
2119/// Options for a copy request
2120#[derive(Debug, Clone, Default)]
2121pub struct CopyOptions {
2122    /// Configure the [`CopyMode`] for this operation
2123    pub mode: CopyMode,
2124    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
2125    /// that need to pass context-specific information (like tracing spans) via trait methods.
2126    ///
2127    /// These extensions are ignored entirely by backends offered through this crate.
2128    ///
2129    /// They are also excluded from [`PartialEq`] and [`Eq`].
2130    pub extensions: Extensions,
2131}
2132
2133impl CopyOptions {
2134    /// Create a new [`CopyOptions`]
2135    pub fn new() -> Self {
2136        Self::default()
2137    }
2138
2139    /// Sets the `mode.
2140    ///
2141    /// See [`CopyOptions::mode`].
2142    #[must_use]
2143    pub fn with_mode(mut self, mode: CopyMode) -> Self {
2144        self.mode = mode;
2145        self
2146    }
2147
2148    /// Sets the `extensions`.
2149    ///
2150    /// See [`CopyOptions::extensions`].
2151    #[must_use]
2152    pub fn with_extensions(mut self, extensions: Extensions) -> Self {
2153        self.extensions = extensions;
2154        self
2155    }
2156}
2157
2158impl PartialEq<Self> for CopyOptions {
2159    fn eq(&self, other: &Self) -> bool {
2160        let Self {
2161            mode,
2162            extensions: _,
2163        } = self;
2164        let Self {
2165            mode: mode_other,
2166            extensions: _,
2167        } = other;
2168
2169        mode == mode_other
2170    }
2171}
2172
2173impl Eq for CopyOptions {}
2174
2175/// Configure preconditions for the target of rename operation.
2176///
2177/// Note though that the source location may or not be deleted at the same time in an atomic operation. There is
2178/// currently NO flag to control the atomicity of "delete source at the same time as creating the target".
2179#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
2180pub enum RenameTargetMode {
2181    /// Perform a write operation on the target, overwriting any object present at the provided path.
2182    #[default]
2183    Overwrite,
2184    /// Perform an atomic write operation of the target, returning [`Error::AlreadyExists`] if an
2185    /// object already exists at the provided path.
2186    Create,
2187}
2188
2189/// Options for a rename request
2190#[derive(Debug, Clone, Default)]
2191pub struct RenameOptions {
2192    /// Configure the [`RenameTargetMode`] for this operation
2193    pub target_mode: RenameTargetMode,
2194    /// Implementation-specific extensions. Intended for use by [`ObjectStore`] implementations
2195    /// that need to pass context-specific information (like tracing spans) via trait methods.
2196    ///
2197    /// These extensions are ignored entirely by backends offered through this crate.
2198    ///
2199    /// They are also excluded from [`PartialEq`] and [`Eq`].
2200    pub extensions: Extensions,
2201}
2202
2203impl RenameOptions {
2204    /// Create a new [`RenameOptions`]
2205    pub fn new() -> Self {
2206        Self::default()
2207    }
2208
2209    /// Sets the `target_mode=.
2210    ///
2211    /// See [`RenameOptions::target_mode`].
2212    #[must_use]
2213    pub fn with_target_mode(mut self, target_mode: RenameTargetMode) -> Self {
2214        self.target_mode = target_mode;
2215        self
2216    }
2217
2218    /// Sets the `extensions`.
2219    ///
2220    /// See [`RenameOptions::extensions`].
2221    #[must_use]
2222    pub fn with_extensions(mut self, extensions: Extensions) -> Self {
2223        self.extensions = extensions;
2224        self
2225    }
2226}
2227
2228impl PartialEq<Self> for RenameOptions {
2229    fn eq(&self, other: &Self) -> bool {
2230        let Self {
2231            target_mode,
2232            extensions: _,
2233        } = self;
2234        let Self {
2235            target_mode: target_mode_other,
2236            extensions: _,
2237        } = other;
2238
2239        target_mode == target_mode_other
2240    }
2241}
2242
2243impl Eq for RenameOptions {}
2244
2245/// A specialized `Result` for object store-related errors
2246pub type Result<T, E = Error> = std::result::Result<T, E>;
2247
2248/// A specialized `Error` for object store-related errors
2249#[derive(Debug, thiserror::Error)]
2250#[non_exhaustive]
2251pub enum Error {
2252    /// A fallback error type when no variant matches
2253    #[error("Generic {} error: {}", store, source)]
2254    Generic {
2255        /// The store this error originated from
2256        store: &'static str,
2257        /// The wrapped error
2258        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2259    },
2260
2261    /// Error when the object is not found at given location
2262    #[error("Object at location {} not found: {}", path, source)]
2263    NotFound {
2264        /// The path to file
2265        path: String,
2266        /// The wrapped error
2267        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2268    },
2269
2270    /// Error for invalid path
2271    #[error("Encountered object with invalid path: {}", source)]
2272    InvalidPath {
2273        /// The wrapped error
2274        #[from]
2275        source: path::Error,
2276    },
2277
2278    /// Error when `tokio::spawn` failed
2279    #[cfg(feature = "tokio")]
2280    #[error("Error joining spawned task: {}", source)]
2281    JoinError {
2282        /// The wrapped error
2283        #[from]
2284        source: tokio::task::JoinError,
2285    },
2286
2287    /// Error when the attempted operation is not supported
2288    #[error("Operation not supported: {}", source)]
2289    NotSupported {
2290        /// The wrapped error
2291        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2292    },
2293
2294    /// Error when the object already exists
2295    #[error("Object at location {} already exists: {}", path, source)]
2296    AlreadyExists {
2297        /// The path to the
2298        path: String,
2299        /// The wrapped error
2300        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2301    },
2302
2303    /// Error when the required conditions failed for the operation
2304    #[error("Request precondition failure for path {}: {}", path, source)]
2305    Precondition {
2306        /// The path to the file
2307        path: String,
2308        /// The wrapped error
2309        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2310    },
2311
2312    /// Error when the object at the location isn't modified
2313    #[error("Object at location {} not modified: {}", path, source)]
2314    NotModified {
2315        /// The path to the file
2316        path: String,
2317        /// The wrapped error
2318        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2319    },
2320
2321    /// Error when an operation is not implemented
2322    #[error("Operation {operation} not yet implemented by {implementer}.")]
2323    NotImplemented {
2324        /// What isn't implemented. Should include at least the method
2325        /// name that was called; could also include other relevant
2326        /// subcontexts.
2327        operation: String,
2328
2329        /// Which driver this is that hasn't implemented this operation,
2330        /// to aid debugging in contexts that may be using multiple implementations.
2331        implementer: String,
2332    },
2333
2334    /// Error when the used credentials don't have enough permission
2335    /// to perform the requested operation
2336    #[error(
2337        "The operation lacked the necessary privileges to complete for path {}: {}",
2338        path,
2339        source
2340    )]
2341    PermissionDenied {
2342        /// The path to the file
2343        path: String,
2344        /// The wrapped error
2345        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2346    },
2347
2348    /// Error when the used credentials lack valid authentication
2349    #[error(
2350        "The operation lacked valid authentication credentials for path {}: {}",
2351        path,
2352        source
2353    )]
2354    Unauthenticated {
2355        /// The path to the file
2356        path: String,
2357        /// The wrapped error
2358        source: Box<dyn std::error::Error + Send + Sync + 'static>,
2359    },
2360
2361    /// Error when a configuration key is invalid for the store used
2362    #[error("Configuration key: '{}' is not valid for store '{}'.", key, store)]
2363    UnknownConfigurationKey {
2364        /// The object store used
2365        store: &'static str,
2366        /// The configuration key used
2367        key: String,
2368    },
2369}
2370
2371impl From<Error> for std::io::Error {
2372    fn from(e: Error) -> Self {
2373        let kind = match &e {
2374            Error::NotFound { .. } => std::io::ErrorKind::NotFound,
2375            _ => std::io::ErrorKind::Other,
2376        };
2377        Self::new(kind, e)
2378    }
2379}
2380
2381#[cfg(test)]
2382mod tests {
2383    use super::*;
2384
2385    use chrono::TimeZone;
2386
2387    macro_rules! maybe_skip_integration {
2388        () => {
2389            if std::env::var("TEST_INTEGRATION").is_err() {
2390                eprintln!("Skipping integration test - set TEST_INTEGRATION");
2391                return;
2392            }
2393        };
2394    }
2395    pub(crate) use maybe_skip_integration;
2396
2397    /// Skip a test that asserts the *storage backend* rejects an invalid presigned request
2398    /// (tampered signature, expired URL, unmet signed header) unless
2399    /// `TEST_S3_SIGNATURE_ENFORCEMENT` is set.
2400    ///
2401    /// These tests require a backend that actually validates SigV4 — real S3 or MinIO. The
2402    /// LocalStack emulator used by the default integration suite does not validate presigned
2403    /// signatures or expiry and would return success, so they are gated separately. See
2404    /// CONTRIBUTING.md for the MinIO/S3 setup.
2405    macro_rules! maybe_skip_signature_enforcement {
2406        () => {
2407            if std::env::var("TEST_S3_SIGNATURE_ENFORCEMENT").is_err() {
2408                eprintln!(
2409                    "Skipping signature-enforcement test - set TEST_S3_SIGNATURE_ENFORCEMENT \
2410                     and point at a backend that validates SigV4 (real S3 or MinIO)"
2411                );
2412                return;
2413            }
2414        };
2415    }
2416    pub(crate) use maybe_skip_signature_enforcement;
2417
2418    /// Test that the returned stream does not borrow the lifetime of Path
2419    fn list_store<'a>(
2420        store: &'a dyn ObjectStore,
2421        path_str: &str,
2422    ) -> BoxStream<'a, Result<ObjectMeta>> {
2423        let path = Path::from(path_str);
2424        store.list(Some(&path))
2425    }
2426
2427    #[cfg(any(feature = "azure-base", feature = "aws-base"))]
2428    pub(crate) async fn signing<T>(integration: &T)
2429    where
2430        T: ObjectStore + signer::Signer,
2431    {
2432        use ::http::Method;
2433        use std::time::Duration;
2434
2435        let data = Bytes::from("hello world");
2436        let path = Path::from("file.txt");
2437        integration.put(&path, data.clone().into()).await.unwrap();
2438
2439        let signed = integration
2440            .signed_url(Method::GET, &path, Duration::from_secs(60))
2441            .await
2442            .unwrap();
2443
2444        let resp = reqwest::get(signed).await.unwrap();
2445        let loaded = resp.bytes().await.unwrap();
2446
2447        assert_eq!(data, loaded);
2448    }
2449
2450    #[cfg(any(feature = "aws-base", feature = "azure-base"))]
2451    pub(crate) async fn tagging<F, Fut>(storage: Arc<dyn ObjectStore>, validate: bool, get_tags: F)
2452    where
2453        F: Fn(Path) -> Fut + Send + Sync,
2454        Fut: std::future::Future<Output = Result<client::HttpResponse>> + Send,
2455    {
2456        use bytes::Buf;
2457        use serde::Deserialize;
2458        use tokio::io::AsyncWriteExt;
2459
2460        use crate::buffered::BufWriter;
2461
2462        #[derive(Deserialize)]
2463        struct Tagging {
2464            #[serde(rename = "TagSet")]
2465            list: TagList,
2466        }
2467
2468        #[derive(Debug, Deserialize)]
2469        struct TagList {
2470            #[serde(rename = "Tag")]
2471            tags: Vec<Tag>,
2472        }
2473
2474        #[derive(Debug, Deserialize, Eq, PartialEq)]
2475        #[serde(rename_all = "PascalCase")]
2476        struct Tag {
2477            key: String,
2478            value: String,
2479        }
2480
2481        let tags = vec![
2482            Tag {
2483                key: "foo.com=bar/s".to_string(),
2484                value: "bananas/foo.com-_".to_string(),
2485            },
2486            Tag {
2487                key: "namespace/key.foo".to_string(),
2488                value: "value with a space".to_string(),
2489            },
2490        ];
2491        let mut tag_set = TagSet::default();
2492        for t in &tags {
2493            tag_set.push(&t.key, &t.value)
2494        }
2495
2496        let path = Path::from("tag_test");
2497        storage
2498            .put_opts(&path, "test".into(), tag_set.clone().into())
2499            .await
2500            .unwrap();
2501
2502        let multi_path = Path::from("tag_test_multi");
2503        let mut write = storage
2504            .put_multipart_opts(&multi_path, tag_set.clone().into())
2505            .await
2506            .unwrap();
2507
2508        write.put_part("foo".into()).await.unwrap();
2509        write.complete().await.unwrap();
2510
2511        let buf_path = Path::from("tag_test_buf");
2512        let mut buf = BufWriter::new(storage, buf_path.clone()).with_tags(tag_set);
2513        buf.write_all(b"foo").await.unwrap();
2514        buf.shutdown().await.unwrap();
2515
2516        // Write should always succeed, but certain configurations may simply ignore tags
2517        if !validate {
2518            return;
2519        }
2520
2521        for path in [path, multi_path, buf_path] {
2522            let resp = get_tags(path.clone()).await.unwrap();
2523            let body = resp.into_body().bytes().await.unwrap();
2524
2525            let mut resp: Tagging = quick_xml::de::from_reader(body.reader()).unwrap();
2526            resp.list.tags.sort_by(|a, b| a.key.cmp(&b.key));
2527            assert_eq!(resp.list.tags, tags);
2528        }
2529    }
2530
2531    #[tokio::test]
2532    async fn test_list_lifetimes() {
2533        let store = memory::InMemory::new();
2534        let mut stream = list_store(&store, "path");
2535        assert!(stream.next().await.is_none());
2536    }
2537
2538    #[test]
2539    fn test_preconditions() {
2540        let mut meta = ObjectMeta {
2541            location: Path::from("test"),
2542            last_modified: Utc.timestamp_nanos(100),
2543            size: 100,
2544            e_tag: Some("\"123\"".to_string()),
2545            version: None,
2546        };
2547
2548        let mut options = GetOptions::default();
2549        options.check_preconditions(&meta).unwrap();
2550
2551        options.if_modified_since = Some(Utc.timestamp_nanos(50));
2552        options.check_preconditions(&meta).unwrap();
2553
2554        options.if_modified_since = Some(Utc.timestamp_nanos(100));
2555        options.check_preconditions(&meta).unwrap_err();
2556
2557        options.if_modified_since = Some(Utc.timestamp_nanos(101));
2558        options.check_preconditions(&meta).unwrap_err();
2559
2560        options = GetOptions::default();
2561
2562        options.if_unmodified_since = Some(Utc.timestamp_nanos(50));
2563        options.check_preconditions(&meta).unwrap_err();
2564
2565        options.if_unmodified_since = Some(Utc.timestamp_nanos(100));
2566        options.check_preconditions(&meta).unwrap();
2567
2568        options.if_unmodified_since = Some(Utc.timestamp_nanos(101));
2569        options.check_preconditions(&meta).unwrap();
2570
2571        options = GetOptions::default();
2572
2573        options.if_match = Some("\"123\"".to_string());
2574        options.check_preconditions(&meta).unwrap();
2575
2576        options.if_match = Some("\"123\",\"354\"".to_string());
2577        options.check_preconditions(&meta).unwrap();
2578
2579        options.if_match = Some("\"354\", \"123\"".to_string());
2580        options.check_preconditions(&meta).unwrap();
2581
2582        options.if_match = Some("\"354\"".to_string());
2583        options.check_preconditions(&meta).unwrap_err();
2584
2585        options.if_match = Some("*".to_string());
2586        options.check_preconditions(&meta).unwrap();
2587
2588        // If-Match takes precedence
2589        options.if_unmodified_since = Some(Utc.timestamp_nanos(200));
2590        options.check_preconditions(&meta).unwrap();
2591
2592        options = GetOptions::default();
2593
2594        options.if_none_match = Some("\"123\"".to_string());
2595        options.check_preconditions(&meta).unwrap_err();
2596
2597        options.if_none_match = Some("*".to_string());
2598        options.check_preconditions(&meta).unwrap_err();
2599
2600        options.if_none_match = Some("\"1232\"".to_string());
2601        options.check_preconditions(&meta).unwrap();
2602
2603        options.if_none_match = Some("\"23\", \"123\"".to_string());
2604        options.check_preconditions(&meta).unwrap_err();
2605
2606        // If-None-Match takes precedence
2607        options.if_modified_since = Some(Utc.timestamp_nanos(10));
2608        options.check_preconditions(&meta).unwrap_err();
2609
2610        // Check missing ETag
2611        meta.e_tag = None;
2612        options = GetOptions::default();
2613
2614        options.if_none_match = Some("*".to_string()); // Fails if any file exists
2615        options.check_preconditions(&meta).unwrap_err();
2616
2617        options = GetOptions::default();
2618        options.if_match = Some("*".to_string()); // Passes if file exists
2619        options.check_preconditions(&meta).unwrap();
2620    }
2621
2622    #[test]
2623    #[cfg(feature = "http-base")]
2624    fn test_reexported_types() {
2625        // Test HeaderMap
2626        let mut headers = HeaderMap::new();
2627        headers.insert("content-type", HeaderValue::from_static("text/plain"));
2628        assert_eq!(headers.len(), 1);
2629
2630        // Test HeaderValue
2631        let value = HeaderValue::from_static("test-value");
2632        assert_eq!(value.as_bytes(), b"test-value");
2633
2634        // Test Extensions
2635        let mut extensions = Extensions::new();
2636        extensions.insert("test-key");
2637        assert!(extensions.get::<&str>().is_some());
2638    }
2639
2640    #[test]
2641    fn test_get_options_builder() {
2642        let dt = Utc::now();
2643        let extensions = Extensions::new();
2644
2645        let options = GetOptions::new();
2646
2647        // assert defaults
2648        assert_eq!(options.if_match, None);
2649        assert_eq!(options.if_none_match, None);
2650        assert_eq!(options.if_modified_since, None);
2651        assert_eq!(options.if_unmodified_since, None);
2652        assert_eq!(options.range, None);
2653        assert_eq!(options.version, None);
2654        assert!(!options.head);
2655        assert!(options.extensions.get::<&str>().is_none());
2656
2657        let options = options
2658            .with_if_match(Some("etag-match"))
2659            .with_if_none_match(Some("etag-none-match"))
2660            .with_if_modified_since(Some(dt))
2661            .with_if_unmodified_since(Some(dt))
2662            .with_range(Some(0..100))
2663            .with_version(Some("version-1"))
2664            .with_head(true)
2665            .with_extensions(extensions.clone());
2666
2667        assert_eq!(options.if_match, Some("etag-match".to_string()));
2668        assert_eq!(options.if_none_match, Some("etag-none-match".to_string()));
2669        assert_eq!(options.if_modified_since, Some(dt));
2670        assert_eq!(options.if_unmodified_since, Some(dt));
2671        assert_eq!(options.range, Some(GetRange::Bounded(0..100)));
2672        assert_eq!(options.version, Some("version-1".to_string()));
2673        assert!(options.head);
2674        assert_eq!(options.extensions.get::<&str>(), extensions.get::<&str>());
2675    }
2676
2677    fn takes_generic_object_store<T: ObjectStore>(store: T) {
2678        // This function is just to ensure that the trait bounds are satisfied
2679        let _ = store;
2680    }
2681    #[test]
2682    fn test_dyn_impl() {
2683        let store: Arc<dyn ObjectStore> = Arc::new(memory::InMemory::new());
2684        takes_generic_object_store(store);
2685        let store: Box<dyn ObjectStore> = Box::new(memory::InMemory::new());
2686        takes_generic_object_store(store);
2687    }
2688    #[test]
2689    fn test_generic_impl() {
2690        let store = Arc::new(memory::InMemory::new());
2691        takes_generic_object_store(store);
2692        let store = Box::new(memory::InMemory::new());
2693        takes_generic_object_store(store);
2694    }
2695}