Skip to main content

lfsx_server/storage/s3/
keyspace.rs

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    // None when the store's timestamp cannot be read, which is treated as "too
31    // young to touch": deleting somebody's upload on the strength of a date this
32    // server could not parse is the wrong way to be wrong.
33    pub(crate) fn age(&self) -> Option<Duration> {
34        Duration::try_from(time::OffsetDateTime::now_utc() - self.modified?).ok()
35    }
36}
37
38// An href a client uses directly, and the headers it has to send with it. The
39// headers are part of the signature, so they are not advice.
40pub struct Presigned {
41    pub href: String,
42    pub headers: Vec<(String, String)>,
43}
44
45// The bucket as a keyspace: whole values written, read, deleted and listed by
46// key, with the signing and the HTTP client in one place. It knows nothing about
47// objects, oids or repositories: what a key means is decided a layer up, which
48// is what lets the lock store share the bucket with the object store without
49// either of them reaching into the other.
50#[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    // Every key under a prefix, following the continuation token to the end.
164    // Stopping at the first page would report a repository holding a thousand
165    // locks as holding a thousand and none of the rest, and a lock nobody can
166    // see is a lock nobody respects.
167    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    // A listing that says whether it finished. Collection needs the difference:
177    // concluding "no marker anywhere references this object" from a listing that
178    // stopped halfway is how a sweep deletes bytes another repository still
179    // holds. Everything else wants the strict form and gets `entries`.
180    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
197// A conditional write that is refused makes the store answer and hang up, and
198// the connection goes back into the pool looking usable. The next request on it
199// fails at the transport layer with nothing to do with the store's health, which
200// is how a losing `git lfs lock` came back as a 500 instead of a 409.
201//
202// Retried once, and only for requests that carry no body: a GET and a HEAD can be
203// repeated with no consequence, so a dead connection costs a round trip rather
204// than an error. A PUT is not retried here.
205pub(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}