1mod 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#[derive(Clone)]
100pub struct S3Client {
101 inner: Arc<ClientInner>,
102}
103
104impl S3Client {
105 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 pub fn config(&self) -> &S3Config {
123 &self.inner.config
124 }
125
126 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}