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;