Skip to main content

lfsx_server/storage/
s3.rs

1pub(crate) mod collect;
2pub(crate) mod keyspace;
3pub(crate) mod multipart;
4pub(crate) mod probe;
5pub(crate) mod refs;
6pub(crate) mod sizes;
7pub(crate) mod usage;
8
9use std::time::Duration;
10
11use axum::body::Bytes;
12use futures_util::Stream;
13
14use base64::Engine;
15
16use crate::error::Error;
17use crate::namespace::Namespace;
18use crate::oid::Oid;
19use crate::storage::Reclaimed;
20
21pub(crate) use keyspace::S3Keys;
22pub use keyspace::{AzureConfig, AzureKeys, GcsConfig, GcsKeys, Keyspace, Presigned};
23
24// Enough to hide the round trips a bucket charges for without becoming a burst
25// the store answers with 503, and the same figure the batch endpoint settled on
26// for the same reason.
27const SIZES_AT_ONCE: usize = 16;
28
29pub struct S3Config {
30    pub endpoint: String,
31    pub bucket: String,
32    pub region: String,
33    pub access_key: String,
34    pub secret_key: String,
35    pub path_style: bool,
36    // How long a signature is good for. It is the same number the batch
37    // response advertises as `expires_in`, because a client told it has half an
38    // hour and handed a URL that dies in five minutes will fail a resume it had
39    // every reason to expect to work.
40    pub lifetime: Duration,
41}
42
43// The same layout as the local store, for the same reasons. The bytes live once
44// under a key derived from their digest, and a repository that holds them owns
45// an empty marker beside it, the object store's answer to a hard link. It is
46// what keeps two projects sharing an asset pack from paying twice, and what
47// stops a repository reading an object it never pushed: the marker is the proof
48// of possession, and it is the only thing the permission check consults.
49//
50// Everything below is object semantics. What it takes to talk to the store at
51// all (signing, retrying, listing, the client) is the keyspace underneath.
52#[derive(Clone)]
53pub struct S3Store {
54    keys: Keyspace,
55    redirect: bool,
56}
57
58impl S3Store {
59    pub fn new(keys: Keyspace, redirect: bool) -> Self {
60        Self { keys, redirect }
61    }
62
63    fn content_key(oid: &Oid) -> String {
64        let (first, second) = oid.fanout();
65        format!(".content/{first}/{second}/{oid}")
66    }
67
68    // Where a client uploads to when the bytes never pass through this server.
69    // Per repository on purpose: the shared content key would take bytes from
70    // anyone allowed to write, and then nothing distinguishes a repository that
71    // uploaded an object from one that merely knew its digest. A key only this
72    // repository was handed a signature for is the proof of possession that the
73    // marker stands for everywhere else.
74    fn incoming_key(ns: &Namespace, oid: &Oid) -> String {
75        let (first, second) = oid.fanout();
76        format!(
77            ".incoming/{}/{}/{first}/{second}/{oid}",
78            ns.stored_org(),
79            ns.repo()
80        )
81    }
82
83    fn marker_key(ns: &Namespace, oid: &Oid) -> String {
84        let (first, second) = oid.fanout();
85        format!("{}/{}/{first}/{second}/{oid}", ns.stored_org(), ns.repo())
86    }
87
88    fn own_prefix(ns: &Namespace) -> String {
89        format!("{}/{}/", ns.stored_org(), ns.repo())
90    }
91
92    pub(crate) async fn read_meta(&self, key: &str) -> Result<Option<Vec<u8>>, Error> {
93        self.keys.get_bytes(key).await
94    }
95
96    pub(crate) async fn write_meta(&self, key: &str, bytes: Vec<u8>) -> Result<(), Error> {
97        self.keys.put(key, bytes).await
98    }
99
100    pub(crate) async fn delete_meta(&self, key: &str) -> Result<(), Error> {
101        self.keys.delete(key).await.map(|_| ())
102    }
103
104    pub async fn reachable(&self) -> Result<(), Error> {
105        self.keys.reachable().await
106    }
107
108    pub async fn exists(&self, ns: &Namespace, oid: &Oid) -> bool {
109        self.keys.head(&Self::marker_key(ns, oid)).await.is_ok()
110    }
111
112    pub async fn size_of(&self, oid: &Oid) -> Result<u64, Error> {
113        self.keys.head(&Self::content_key(oid)).await
114    }
115
116    // A download is streamed through this server rather than redirected, so the
117    // features that live in the byte path (the counters, the ranges, and the
118    // compression that will follow) keep working. The pre-signed redirect is a
119    // separate mode for operators who would rather spend the object store's
120    // bandwidth than their own.
121    pub async fn read(
122        &self,
123        oid: &Oid,
124        start: u64,
125        length: u64,
126    ) -> Result<impl Stream<Item = Result<Bytes, reqwest::Error>> + use<>, Error> {
127        self.keys
128            .get_range(&Self::content_key(oid), start, length)
129            .await
130    }
131
132    // A URL the client fetches from the bucket directly, so the bytes never
133    // cross this server. Whether the caller is entitled to them has already been
134    // settled by the marker before this is called: the signature is scoped to
135    // one content key and expires, and it grants nothing the batch response was
136    // not about to grant anyway.
137    pub fn presigned_download(&self, oid: &Oid) -> Option<String> {
138        if !self.redirect {
139            return None;
140        }
141
142        self.keys.signed_download(&Self::content_key(oid))
143    }
144
145    // A URL the client PUTs the object to, and the headers it has to send with
146    // it. The digest is bound into the signature, so the store refuses anything
147    // that does not hash to the object it was signed for: a client with this URL
148    // cannot put arbitrary bytes anywhere, which is what makes handing one out
149    // safe at all.
150    //
151    // None above the single-request ceiling, and that is not a refusal: the
152    // object falls back to coming through this server, which sends it in parts.
153    // A client cannot do the same, because the `basic` transfer adapter every
154    // git-lfs speaks does one PUT to one href and has nowhere to put a second.
155    // So the ceiling multipart removes for the streamed path is real and
156    // permanent for this one, and the only question is whether the client learns
157    // it now or after uploading five gigabytes.
158    //
159    // It also keeps `adopt` honest: `CopyObject` stops at the same 5 GiB, and
160    // nothing can reach `.incoming/` above it while this holds.
161    pub fn presigned_upload(&self, ns: &Namespace, oid: &Oid, size: u64) -> Option<Presigned> {
162        if !self.redirect || size > multipart::SINGLE_PUT_CEILING {
163            return None;
164        }
165
166        let digest =
167            base64::engine::general_purpose::STANDARD.encode(hex::decode(oid.as_str()).ok()?);
168
169        self.keys
170            .signed_upload(&Self::incoming_key(ns, oid), &digest)
171    }
172
173    // How big the object a client uploaded actually is, which is the first thing
174    // this server learns about it: nothing measured the bytes on the way past.
175    pub async fn uploaded_size(&self, ns: &Namespace, oid: &Oid) -> Result<u64, Error> {
176        self.keys.head(&Self::incoming_key(ns, oid)).await
177    }
178
179    // Take an upload that landed under this repository's own key into the shared
180    // keyspace. The bytes are already known to hash to the oid, because the store
181    // refused everything else.
182    pub async fn adopt(&self, ns: &Namespace, oid: &Oid, size: u64) -> Result<(), Error> {
183        let incoming = Self::incoming_key(ns, oid);
184        let content = Self::content_key(oid);
185
186        // First, before anything here so much as looks at the content.
187        //
188        // The marker is the claim and the ref is the index of it, so a crash
189        // between the two has to leave a ref nobody claims rather than a claim
190        // nothing indexes: the first leaks an object, the second lets a later
191        // sweep free bytes this repository holds.
192        //
193        // Writing it up here rather than beside the marker costs nothing and buys
194        // the race below. A sweep asks the index one last time before deleting
195        // bytes, so a claim recorded before this repository even checked whether
196        // the content exists is a claim that sweep will see.
197        refs::write(&self.keys, ns, oid).await?;
198
199        // Already there means another repository pushed the same object, and the
200        // bytes are identical by construction.
201        if self.keys.head(&content).await.is_err() {
202            self.keys.copy(&incoming, &content).await?;
203        }
204
205        self.keys
206            .put(&Self::marker_key(ns, oid), Vec::new())
207            .await?;
208
209        sizes::write(&self.keys, ns, oid, size).await?;
210
211        // Leaving it would pay for the object twice. A failure here is not worth
212        // failing the push over: the object is adopted, and what is left is a key
213        // the operator can see.
214        if let Err(error) = self.keys.delete(&incoming).await {
215            tracing::warn!(%error, key = incoming, "an adopted upload could not be cleaned up");
216        }
217
218        Ok(())
219    }
220
221    // The upload has already been streamed to a staging file, hashed and checked
222    // against everything the server enforces, so that file is what goes up,
223    // streamed from disk rather than read into memory, because an object here is
224    // measured in gigabytes and the whole storage layer is built on holding at
225    // most a few megabytes of one at a time.
226    //
227    // The bytes go up once, keyed by their digest, and the marker records that
228    // this repository holds them. Content that is already there is skipped: the
229    // key would receive the same bytes it already has.
230    pub async fn store(
231        &self,
232        ns: &Namespace,
233        oid: &Oid,
234        staged: &std::path::Path,
235    ) -> Result<(), Error> {
236        // Before the content is even looked at, for the reason `adopt` gives:
237        // this is what a sweep re-reads before deleting bytes, so a claim
238        // recorded here cannot be missed by one that is already deciding.
239        refs::write(&self.keys, ns, oid).await?;
240
241        // Read here rather than inside the branch below, because the size index
242        // wants it whether or not these bytes are the ones that go up: an object
243        // another repository pushed first is still this repository's to account
244        // for.
245        let length = tokio::fs::metadata(staged).await?.len();
246
247        if self.keys.head(&Self::content_key(oid)).await.is_err() {
248            self.keys.put_file(&Self::content_key(oid), staged).await?;
249        }
250
251        self.keys
252            .put(&Self::marker_key(ns, oid), Vec::new())
253            .await?;
254
255        sizes::write(&self.keys, ns, oid, length).await
256    }
257
258    // What an interrupted upload leaves behind. A client can negotiate, PUT the
259    // object, and never report it: the bytes sit under its own upload key and
260    // nothing else will ever look at them. The local path has had a reclaimer for
261    // this since the beginning, and a bucket had none, so the cost was unbounded
262    // over time and invisible.
263    pub async fn reclaim_incoming(&self, older_than: Duration) -> Result<Reclaimed, Error> {
264        let mut reclaimed = Reclaimed::default();
265
266        // `.probe/` too. A startup probe draws a key nothing else uses so that no
267        // run can read another's leftovers, which means a run that dies before
268        // cleaning up leaves one behind rather than overwriting it. They are
269        // empty or nearly so, and this is already the sweep for writes nobody
270        // will ever come back for.
271        for prefix in [".incoming/", ".probe/"] {
272            for entry in self.keys.entries(prefix).await? {
273                // A slow client on a bad connection is not an abandoned one.
274                if entry.age().is_none_or(|age| age < older_than) {
275                    continue;
276                }
277
278                if self.keys.delete(&entry.key).await.is_ok() {
279                    reclaimed.files += 1;
280                    reclaimed.bytes += entry.size;
281                }
282            }
283        }
284
285        Ok(reclaimed)
286    }
287}
288
289#[cfg(test)]
290pub(crate) mod tests;