Skip to main content

lfsx_server/storage/s3/
keyspace.rs

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