lfsx_server/storage/s3/
keyspace.rs1mod s3;
2
3use std::path::Path;
4use std::time::Duration;
5
6use axum::body::Bytes;
7use futures_util::Stream;
8
9use crate::error::Error;
10
11pub(crate) use s3::S3Keys;
12
13pub(crate) struct Listing {
14 pub(crate) entries: Vec<Entry>,
15 pub(crate) complete: bool,
16}
17
18pub(crate) struct Entry {
19 pub(crate) key: String,
20 pub(crate) modified: Option<time::OffsetDateTime>,
21 pub(crate) size: u64,
22}
23
24impl Entry {
25 pub(crate) fn age(&self) -> Option<Duration> {
29 Duration::try_from(time::OffsetDateTime::now_utc() - self.modified?).ok()
30 }
31}
32
33pub struct Presigned {
36 pub href: String,
37 pub headers: Vec<(String, String)>,
38}
39
40#[derive(Clone)]
46pub enum Keyspace {
47 S3(S3Keys),
48}
49
50impl Keyspace {
51 pub(crate) async fn reachable(&self) -> Result<(), Error> {
52 match self {
53 Self::S3(keys) => keys.reachable().await,
54 }
55 }
56
57 pub(crate) fn signed_download(&self, key: &str) -> Option<String> {
58 match self {
59 Self::S3(keys) => Some(keys.signed_download(key)),
60 }
61 }
62
63 pub(crate) fn signed_upload(&self, key: &str, digest: &str) -> Option<Presigned> {
64 match self {
65 Self::S3(keys) => Some(keys.signed_upload(key, digest)),
66 }
67 }
68
69 pub(crate) async fn get_range(
70 &self,
71 key: &str,
72 start: u64,
73 length: u64,
74 ) -> Result<impl Stream<Item = Result<Bytes, reqwest::Error>> + use<>, Error> {
75 let response = match self {
76 Self::S3(keys) => keys.get_range(key, start, length).await?,
77 };
78
79 Ok(response.bytes_stream())
80 }
81
82 pub(crate) async fn head(&self, key: &str) -> Result<u64, Error> {
83 match self {
84 Self::S3(keys) => keys.head(key).await,
85 }
86 }
87
88 pub(crate) async fn put(&self, key: &str, body: Vec<u8>) -> Result<(), Error> {
89 let length = body.len() as u64;
90
91 match self {
92 Self::S3(keys) => keys.put(key, reqwest::Body::from(body), length).await,
93 }
94 }
95
96 pub(crate) async fn put_file(&self, key: &str, staged: &Path) -> Result<(), Error> {
97 match self {
98 Self::S3(keys) => keys.put_file(key, staged).await,
99 }
100 }
101
102 pub(crate) async fn copy(&self, from: &str, to: &str) -> Result<(), Error> {
103 match self {
104 Self::S3(keys) => keys.copy(from, to).await,
105 }
106 }
107
108 pub(crate) async fn put_if_absent(&self, key: &str, body: Vec<u8>) -> Result<bool, Error> {
109 match self {
110 Self::S3(keys) => keys.put_if_absent(key, body).await,
111 }
112 }
113
114 pub(crate) async fn get_bytes(&self, key: &str) -> Result<Option<Vec<u8>>, Error> {
115 match self {
116 Self::S3(keys) => keys.get_bytes(key).await,
117 }
118 }
119
120 pub(crate) async fn delete(&self, key: &str) -> Result<bool, Error> {
121 match self {
122 Self::S3(keys) => keys.delete(key).await,
123 }
124 }
125
126 pub(crate) async fn entries(&self, prefix: &str) -> Result<Vec<Entry>, Error> {
127 match self {
128 Self::S3(keys) => keys.entries(prefix).await,
129 }
130 }
131
132 pub(crate) async fn keys(&self, prefix: &str) -> Result<Vec<String>, Error> {
137 Ok(self
138 .entries(prefix)
139 .await?
140 .into_iter()
141 .map(|entry| entry.key)
142 .collect())
143 }
144
145 pub(crate) async fn listing(&self, prefix: &str) -> Listing {
150 match self.entries(prefix).await {
151 Ok(entries) => Listing {
152 entries,
153 complete: true,
154 },
155 Err(error) => {
156 tracing::warn!(%error, prefix, "the listing could not be finished");
157 Listing {
158 entries: Vec::new(),
159 complete: false,
160 }
161 }
162 }
163 }
164}
165
166pub(crate) async fn read_retrying(
175 request: reqwest::RequestBuilder,
176) -> Result<reqwest::Response, Error> {
177 let retry = request.try_clone();
178
179 match request.send().await {
180 Ok(response) => Ok(response),
181 Err(_) => match retry {
182 Some(retry) => retry.send().await.map_err(|_| unreachable_store()),
183 None => Err(unreachable_store()),
184 },
185 }
186}
187
188pub(crate) fn unreachable_store() -> Error {
189 Error::Storage(std::io::Error::other("the object store is unreachable"))
190}
191
192pub(crate) async fn expect_success(
193 response: reqwest::Response,
194 what: &str,
195) -> Result<reqwest::Response, Error> {
196 let status = response.status();
197 if status.is_success() {
198 return Ok(response);
199 }
200
201 let detail = response.text().await.unwrap_or_default();
202
203 Err(Error::Storage(std::io::Error::other(format!(
204 "the object store refused a {what} with {status}: {}",
205 detail.trim()
206 ))))
207}
208
209pub(crate) fn content_length(response: &reqwest::Response) -> Result<u64, Error> {
210 response
211 .headers()
212 .get(reqwest::header::CONTENT_LENGTH)
213 .and_then(|value| value.to_str().ok())
214 .and_then(|value| value.parse().ok())
215 .ok_or_else(|| {
216 Error::Storage(std::io::Error::other(
217 "the object store gave no object size",
218 ))
219 })
220}