Skip to main content

s3_wire/client/
mod.rs

1//! S3 client implementation.
2
3mod download;
4mod multipart;
5mod multipart_upload;
6mod object;
7mod presign;
8mod request;
9
10use std::fmt;
11use std::sync::Arc;
12use std::sync::atomic::{AtomicU32, Ordering};
13
14use crate::config::S3Config;
15use crate::error::S3Error;
16use crate::transport::Transport;
17use crate::{RequestEvent, RequestObserver};
18
19pub use download::DownloadToPathOutput;
20
21struct ClientInner {
22    config: S3Config,
23    transport: Transport,
24    retry_quota: Arc<RetryQuota>,
25}
26
27const RETRY_QUOTA_CAPACITY: u32 = 500;
28const TRANSIENT_RETRY_COST: u32 = 14;
29const THROTTLING_RETRY_COST: u32 = 5;
30
31struct RetryQuota {
32    tokens: AtomicU32,
33}
34
35struct RetryPermit {
36    quota: Arc<RetryQuota>,
37    cost: u32,
38    committed: bool,
39}
40
41impl RetryQuota {
42    fn new() -> Self {
43        Self {
44            tokens: AtomicU32::new(RETRY_QUOTA_CAPACITY),
45        }
46    }
47
48    fn acquire(
49        self: &Arc<Self>,
50        classification: crate::RetryClassification,
51    ) -> Option<RetryPermit> {
52        let cost = match classification {
53            crate::RetryClassification::Throttled => THROTTLING_RETRY_COST,
54            crate::RetryClassification::Retryable => TRANSIENT_RETRY_COST,
55            crate::RetryClassification::Never => return None,
56        };
57        self.tokens
58            .fetch_update(Ordering::AcqRel, Ordering::Acquire, |available| {
59                if available >= cost {
60                    Some(available - cost)
61                } else {
62                    None
63                }
64            })
65            .ok()
66            .map(|_| RetryPermit {
67                quota: Arc::clone(self),
68                cost,
69                committed: false,
70            })
71    }
72
73    fn replenish_after_success(&self, last_retry_cost: Option<u32>) {
74        let restored = last_retry_cost.unwrap_or(1);
75        let _ = self
76            .tokens
77            .fetch_update(Ordering::AcqRel, Ordering::Acquire, |available| {
78                Some(available.saturating_add(restored).min(RETRY_QUOTA_CAPACITY))
79            });
80    }
81}
82
83impl RetryPermit {
84    fn commit(mut self) -> u32 {
85        self.committed = true;
86        self.cost
87    }
88}
89
90impl Drop for RetryPermit {
91    fn drop(&mut self) {
92        if !self.committed {
93            self.quota.replenish_after_success(Some(self.cost));
94        }
95    }
96}
97
98/// An async S3-compatible client.
99#[derive(Clone)]
100pub struct S3Client {
101    inner: Arc<ClientInner>,
102}
103
104impl S3Client {
105    /// Constructs a client from validated configuration.
106    pub fn new(config: S3Config) -> Result<Self, S3Error> {
107        let transport = Transport::new(
108            config.connect_timeout(),
109            config.user_agent(),
110            !config.endpoint().is_https(),
111        )?;
112        Ok(Self {
113            inner: Arc::new(ClientInner {
114                config,
115                transport,
116                retry_quota: Arc::new(RetryQuota::new()),
117            }),
118        })
119    }
120
121    /// Returns this client's validated configuration.
122    pub fn config(&self) -> &S3Config {
123        &self.inner.config
124    }
125
126    /// Creates a cheap bucket-scoped client that shares this client's HTTP
127    /// connection pool and credential provider.
128    ///
129    /// # Errors
130    ///
131    /// Returns an error when `bucket` is invalid for the configured addressing
132    /// style.
133    pub fn for_bucket(&self, bucket: impl Into<String>) -> Result<Self, S3Error> {
134        let config = self.inner.config.for_bucket(bucket)?;
135        Ok(Self {
136            inner: Arc::new(ClientInner {
137                config,
138                transport: self.inner.transport.clone(),
139                retry_quota: Arc::clone(&self.inner.retry_quota),
140            }),
141        })
142    }
143
144    pub(in crate::client) fn observe(&self, event: RequestEvent<'_>) {
145        let Some(observer) = self.inner.config.observer() else {
146            return;
147        };
148        let observer: &dyn RequestObserver = observer.as_ref();
149        let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
150            observer.on_event(event);
151        }));
152    }
153
154    fn operation_target(&self, key: Option<&str>) -> Result<crate::endpoint::EndpointUrl, S3Error> {
155        self.inner.config.endpoint().object_url(
156            self.inner.config.bucket(),
157            key,
158            self.inner.config.addressing_style(),
159        )
160    }
161}
162
163impl fmt::Debug for S3Client {
164    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
165        formatter
166            .debug_struct("S3Client")
167            .field("config", &self.inner.config)
168            .finish_non_exhaustive()
169    }
170}
171
172#[cfg(test)]
173mod tests {
174    use super::*;
175
176    #[test]
177    fn bucket_handles_reuse_validated_service_configuration() {
178        let client =
179            S3Client::new(S3Config::builder().bucket("first-bucket").build().unwrap()).unwrap();
180        let second = client.for_bucket("second-bucket").unwrap();
181        assert_eq!(client.config().bucket(), "first-bucket");
182        assert_eq!(second.config().bucket(), "second-bucket");
183        assert_eq!(client.config().endpoint(), second.config().endpoint());
184        assert!(client.for_bucket("").is_err());
185    }
186
187    #[test]
188    fn retry_quota_is_shared_bounded_and_replenished() {
189        let quota = Arc::new(RetryQuota::new());
190        for _ in 0..35 {
191            let permit = quota
192                .acquire(crate::RetryClassification::Retryable)
193                .expect("quota remains");
194            assert_eq!(permit.commit(), TRANSIENT_RETRY_COST);
195        }
196        assert!(
197            quota
198                .acquire(crate::RetryClassification::Retryable)
199                .is_none()
200        );
201        quota.replenish_after_success(Some(TRANSIENT_RETRY_COST));
202        let permit = quota
203            .acquire(crate::RetryClassification::Retryable)
204            .expect("replenished quota");
205        assert_eq!(permit.commit(), TRANSIENT_RETRY_COST);
206    }
207
208    #[test]
209    fn unused_retry_permit_returns_its_tokens() {
210        let quota = Arc::new(RetryQuota::new());
211        let permit = quota
212            .acquire(crate::RetryClassification::Retryable)
213            .expect("initial quota");
214        drop(permit);
215        assert_eq!(quota.tokens.load(Ordering::Acquire), RETRY_QUOTA_CAPACITY);
216    }
217}