1pub(crate) mod azure;
2pub(crate) mod gcs;
3mod s3;
4mod token;
5
6use std::path::Path;
7use std::time::Duration;
8
9use axum::body::Bytes;
10use futures_util::Stream;
11
12use crate::error::Error;
13
14pub use azure::{AzureConfig, AzureKeys};
15pub use gcs::{GcsConfig, GcsKeys};
16pub(crate) use s3::S3Keys;
17
18pub(crate) struct Listing {
19 pub(crate) entries: Vec<Entry>,
20 pub(crate) complete: bool,
21}
22
23pub(crate) struct Entry {
24 pub(crate) key: String,
25 pub(crate) modified: Option<time::OffsetDateTime>,
26 pub(crate) size: u64,
27}
28
29impl Entry {
30 pub(crate) fn age(&self) -> Option<Duration> {
34 Duration::try_from(time::OffsetDateTime::now_utc() - self.modified?).ok()
35 }
36}
37
38pub struct Presigned {
41 pub href: String,
42 pub headers: Vec<(String, String)>,
43}
44
45#[derive(Clone)]
51pub enum Keyspace {
52 S3(S3Keys),
53 Azure(AzureKeys),
54 Gcs(GcsKeys),
55}
56
57impl Keyspace {
58 pub(crate) async fn reachable(&self) -> Result<(), Error> {
59 match self {
60 Self::S3(keys) => keys.reachable().await,
61 Self::Azure(keys) => keys.reachable().await,
62 Self::Gcs(keys) => keys.reachable().await,
63 }
64 }
65
66 pub(crate) fn signed_download(&self, key: &str) -> Option<String> {
67 match self {
68 Self::S3(keys) => Some(keys.signed_download(key)),
69 Self::Azure(keys) => keys.signed_download(key),
70 Self::Gcs(keys) => keys.signed_download(key),
71 }
72 }
73
74 pub(crate) fn signed_upload(&self, key: &str, digest: &str) -> Option<Presigned> {
75 match self {
76 Self::S3(keys) => Some(keys.signed_upload(key, digest)),
77 Self::Azure(_) | Self::Gcs(_) => None,
78 }
79 }
80
81 pub(crate) async fn get_range(
82 &self,
83 key: &str,
84 start: u64,
85 length: u64,
86 ) -> Result<impl Stream<Item = Result<Bytes, reqwest::Error>> + use<>, Error> {
87 let response = match self {
88 Self::S3(keys) => keys.get_range(key, start, length).await?,
89 Self::Azure(keys) => keys.get_range(key, start, length).await?,
90 Self::Gcs(keys) => keys.get_range(key, start, length).await?,
91 };
92
93 Ok(response.bytes_stream())
94 }
95
96 pub(crate) async fn head(&self, key: &str) -> Result<u64, Error> {
97 match self {
98 Self::S3(keys) => keys.head(key).await,
99 Self::Azure(keys) => keys.head(key).await,
100 Self::Gcs(keys) => keys.head(key).await,
101 }
102 }
103
104 pub(crate) async fn put(&self, key: &str, body: Vec<u8>) -> Result<(), Error> {
105 let length = body.len() as u64;
106
107 match self {
108 Self::S3(keys) => keys.put(key, reqwest::Body::from(body), length).await,
109 Self::Azure(keys) => keys.put(key, reqwest::Body::from(body), length).await,
110 Self::Gcs(keys) => keys.put(key, reqwest::Body::from(body), length).await,
111 }
112 }
113
114 pub(crate) async fn put_file(&self, key: &str, staged: &Path) -> Result<(), Error> {
115 match self {
116 Self::S3(keys) => keys.put_file(key, staged).await,
117 Self::Azure(keys) => keys.put_file(key, staged).await,
118 Self::Gcs(keys) => keys.put_file(key, staged).await,
119 }
120 }
121
122 pub(crate) async fn copy(&self, from: &str, to: &str) -> Result<(), Error> {
123 match self {
124 Self::S3(keys) => keys.copy(from, to).await,
125 Self::Azure(_) | Self::Gcs(_) => Err(Error::Storage(std::io::Error::other(
126 "this object store takes no client uploads, so it has nothing to copy",
127 ))),
128 }
129 }
130
131 pub(crate) async fn put_if_absent(&self, key: &str, body: Vec<u8>) -> Result<bool, Error> {
132 match self {
133 Self::S3(keys) => keys.put_if_absent(key, body).await,
134 Self::Azure(keys) => keys.put_if_absent(key, body).await,
135 Self::Gcs(keys) => keys.put_if_absent(key, body).await,
136 }
137 }
138
139 pub(crate) async fn get_bytes(&self, key: &str) -> Result<Option<Vec<u8>>, Error> {
140 match self {
141 Self::S3(keys) => keys.get_bytes(key).await,
142 Self::Azure(keys) => keys.get_bytes(key).await,
143 Self::Gcs(keys) => keys.get_bytes(key).await,
144 }
145 }
146
147 pub(crate) async fn delete(&self, key: &str) -> Result<bool, Error> {
148 match self {
149 Self::S3(keys) => keys.delete(key).await,
150 Self::Azure(keys) => keys.delete(key).await,
151 Self::Gcs(keys) => keys.delete(key).await,
152 }
153 }
154
155 pub(crate) async fn entries(&self, prefix: &str) -> Result<Vec<Entry>, Error> {
156 match self {
157 Self::S3(keys) => keys.entries(prefix).await,
158 Self::Azure(keys) => keys.entries(prefix).await,
159 Self::Gcs(keys) => keys.entries(prefix).await,
160 }
161 }
162
163 pub(crate) async fn keys(&self, prefix: &str) -> Result<Vec<String>, Error> {
168 Ok(self
169 .entries(prefix)
170 .await?
171 .into_iter()
172 .map(|entry| entry.key)
173 .collect())
174 }
175
176 pub(crate) async fn listing(&self, prefix: &str) -> Listing {
181 match self.entries(prefix).await {
182 Ok(entries) => Listing {
183 entries,
184 complete: true,
185 },
186 Err(error) => {
187 tracing::warn!(%error, prefix, "the listing could not be finished");
188 Listing {
189 entries: Vec::new(),
190 complete: false,
191 }
192 }
193 }
194 }
195}
196
197pub(crate) async fn read_retrying(
206 request: reqwest::RequestBuilder,
207) -> Result<reqwest::Response, Error> {
208 let retry = request.try_clone();
209
210 match request.send().await {
211 Ok(response) => Ok(response),
212 Err(_) => match retry {
213 Some(retry) => retry.send().await.map_err(|_| unreachable_store()),
214 None => Err(unreachable_store()),
215 },
216 }
217}
218
219pub(crate) fn unreachable_store() -> Error {
220 Error::Storage(std::io::Error::other("the object store is unreachable"))
221}
222
223pub(crate) async fn expect_success(
224 response: reqwest::Response,
225 what: &str,
226) -> Result<reqwest::Response, Error> {
227 let status = response.status();
228 if status.is_success() {
229 return Ok(response);
230 }
231
232 let detail = response.text().await.unwrap_or_default();
233
234 Err(Error::Storage(std::io::Error::other(format!(
235 "the object store refused a {what} with {status}: {}",
236 detail.trim()
237 ))))
238}
239
240pub(crate) fn content_length(response: &reqwest::Response) -> Result<u64, Error> {
241 response
242 .headers()
243 .get(reqwest::header::CONTENT_LENGTH)
244 .and_then(|value| value.to_str().ok())
245 .and_then(|value| value.parse().ok())
246 .ok_or_else(|| {
247 Error::Storage(std::io::Error::other(
248 "the object store gave no object size",
249 ))
250 })
251}