Skip to main content

lfsx_server/storage/s3/keyspace/
gcs.rs

1mod credential;
2mod signing;
3
4use std::path::Path;
5use std::sync::Arc;
6use std::time::Duration;
7
8use reqwest::{Method, StatusCode};
9use tokio::io::{AsyncReadExt, AsyncSeekExt};
10
11use super::{Entry, expect_success, read_retrying, unreachable_store};
12use crate::config::GcsCredential;
13use crate::error::Error;
14
15pub(crate) use credential::Credential;
16
17const SINGLE_PUT_CEILING: u64 = 256 * 1024 * 1024;
18const CHUNK: u64 = 64 * 1024 * 1024;
19const RESUME_INCOMPLETE: u16 = 308;
20
21pub struct GcsConfig {
22    pub endpoint: String,
23    pub bucket: String,
24    pub credential: GcsCredential,
25    pub lifetime: Duration,
26}
27
28#[derive(Clone)]
29pub struct GcsKeys {
30    endpoint: reqwest::Url,
31    bucket: String,
32    credential: Arc<Credential>,
33    client: reqwest::Client,
34    lifetime: Duration,
35}
36
37#[derive(serde::Deserialize)]
38struct Object {
39    name: String,
40    size: String,
41    updated: Option<String>,
42}
43
44#[derive(serde::Deserialize)]
45#[serde(rename_all = "camelCase")]
46struct Page {
47    #[serde(default)]
48    items: Vec<Object>,
49    next_page_token: Option<String>,
50}
51
52fn malformed(what: &'static str) -> Error {
53    Error::Storage(std::io::Error::other(what))
54}
55
56impl GcsKeys {
57    pub fn new(config: &GcsConfig) -> Result<Self, Error> {
58        crate::tls::install_crypto_provider();
59
60        Self::with_credential(config, Credential::new(&config.credential)?)
61    }
62
63    pub(crate) fn with_credential(
64        config: &GcsConfig,
65        credential: Credential,
66    ) -> Result<Self, Error> {
67        let mut endpoint = reqwest::Url::parse(&config.endpoint)
68            .map_err(|_| Error::Misconfigured("LFSX_GCS_ENDPOINT is not a URL"))?;
69        if endpoint.cannot_be_a_base() {
70            return Err(Error::Misconfigured("LFSX_GCS_ENDPOINT is not a URL"));
71        }
72        if let Ok(mut segments) = endpoint.path_segments_mut() {
73            segments.pop_if_empty();
74        }
75
76        Ok(Self {
77            endpoint,
78            bucket: config.bucket.clone(),
79            credential: Arc::new(credential),
80            client: reqwest::Client::new(),
81            lifetime: config.lifetime,
82        })
83    }
84
85    fn url(&self, upload: bool, object: Option<&str>, query: &[(&str, &str)]) -> reqwest::Url {
86        let mut url = self.endpoint.clone();
87        if let Ok(mut segments) = url.path_segments_mut() {
88            if upload {
89                segments.push("upload");
90            }
91            segments.extend(["storage", "v1", "b", &self.bucket, "o"]);
92            if let Some(object) = object {
93                segments.push(object);
94            }
95        }
96        if !query.is_empty() {
97            url.query_pairs_mut().extend_pairs(query);
98        }
99        url
100    }
101
102    async fn authorized(
103        &self,
104        method: Method,
105        url: reqwest::Url,
106    ) -> Result<reqwest::RequestBuilder, Error> {
107        let request = self.client.request(method, url);
108
109        Ok(match self.credential.bearer(&self.client).await? {
110            Some(token) => request.bearer_auth(token),
111            None => request,
112        })
113    }
114
115    pub(crate) async fn reachable(&self) -> Result<(), Error> {
116        let url = self.url(false, None, &[("maxResults", "1")]);
117        let response = read_retrying(self.authorized(Method::GET, url).await?).await?;
118
119        if !response.status().is_success() {
120            return Err(Error::Storage(std::io::Error::other(format!(
121                "the object store answered {} for the bucket",
122                response.status()
123            ))));
124        }
125
126        Ok(())
127    }
128
129    pub(crate) fn signed_download(&self, key: &str) -> Option<String> {
130        signing::signed_download(
131            &self.endpoint,
132            &self.bucket,
133            key,
134            self.credential.service_account()?,
135            self.lifetime,
136        )
137    }
138
139    pub(crate) async fn get_range(
140        &self,
141        key: &str,
142        start: u64,
143        length: u64,
144    ) -> Result<reqwest::Response, Error> {
145        let url = self.url(false, Some(key), &[("alt", "media")]);
146        let response = self
147            .authorized(Method::GET, url)
148            .await?
149            .header(
150                reqwest::header::RANGE,
151                format!("bytes={start}-{}", start + length.saturating_sub(1)),
152            )
153            .send()
154            .await
155            .map_err(|_| unreachable_store())?;
156
157        if !response.status().is_success() {
158            return Err(Error::NotFound);
159        }
160
161        Ok(response)
162    }
163
164    pub(crate) async fn head(&self, key: &str) -> Result<u64, Error> {
165        let url = self.url(false, Some(key), &[]);
166        let response = read_retrying(self.authorized(Method::GET, url).await?).await?;
167
168        if !response.status().is_success() {
169            return Err(Error::NotFound);
170        }
171
172        let object: Object = response.json().await.map_err(|_| {
173            malformed("the object store described an object this server could not read")
174        })?;
175
176        object
177            .size
178            .parse()
179            .map_err(|_| malformed("the object store gave no object size"))
180    }
181
182    async fn insert(
183        &self,
184        key: &str,
185        body: reqwest::Body,
186        length: u64,
187        only_if_absent: bool,
188    ) -> Result<reqwest::Response, Error> {
189        let mut query = vec![("uploadType", "media"), ("name", key)];
190        if only_if_absent {
191            query.push(("ifGenerationMatch", "0"));
192        }
193
194        self.authorized(Method::POST, self.url(true, None, &query))
195            .await?
196            .header(reqwest::header::CONTENT_LENGTH, length)
197            .header(reqwest::header::CONTENT_TYPE, "application/octet-stream")
198            .body(body)
199            .send()
200            .await
201            .map_err(|_| unreachable_store())
202    }
203
204    pub(crate) async fn put(
205        &self,
206        key: &str,
207        body: reqwest::Body,
208        length: u64,
209    ) -> Result<(), Error> {
210        let response = self.insert(key, body, length, false).await?;
211        expect_success(response, "write").await?;
212
213        Ok(())
214    }
215
216    pub(crate) async fn put_file(&self, key: &str, staged: &Path) -> Result<(), Error> {
217        let file = tokio::fs::File::open(staged).await?;
218        let length = file.metadata().await?.len();
219
220        if length > SINGLE_PUT_CEILING {
221            drop(file);
222            return self.put_resumable(key, staged, length, CHUNK).await;
223        }
224
225        let stream = tokio_util::io::ReaderStream::new(file);
226        self.put(key, reqwest::Body::wrap_stream(stream), length)
227            .await
228    }
229
230    pub(crate) async fn put_resumable(
231        &self,
232        key: &str,
233        staged: &Path,
234        length: u64,
235        chunk: u64,
236    ) -> Result<(), Error> {
237        let url = self.url(true, None, &[("uploadType", "resumable"), ("name", key)]);
238        let response = self
239            .authorized(Method::POST, url)
240            .await?
241            .header(reqwest::header::CONTENT_LENGTH, 0)
242            .header("x-upload-content-length", length)
243            .send()
244            .await
245            .map_err(|_| unreachable_store())?;
246        let session = expect_success(response, "start a resumable upload")
247            .await?
248            .headers()
249            .get(reqwest::header::LOCATION)
250            .and_then(|value| value.to_str().ok())
251            .map(str::to_owned)
252            .ok_or_else(|| malformed("the object store named no upload session"))?;
253
254        tracing::info!(
255            key,
256            length,
257            chunk,
258            "an object over the single-request ceiling is going up in chunks"
259        );
260
261        let mut offset = 0;
262        while offset < length {
263            let this = chunk.min(length - offset);
264
265            let mut file = tokio::fs::File::open(staged).await?;
266            file.seek(std::io::SeekFrom::Start(offset)).await?;
267            let stream = tokio_util::io::ReaderStream::new(file.take(this));
268
269            let session = self
270                .endpoint
271                .join(&session)
272                .map_err(|_| malformed("the object store named an unusable upload session"))?;
273            let response = self
274                .authorized(Method::PUT, session)
275                .await?
276                .header(reqwest::header::CONTENT_LENGTH, this)
277                .header(
278                    reqwest::header::CONTENT_RANGE,
279                    format!("bytes {offset}-{}/{length}", offset + this - 1),
280                )
281                .body(reqwest::Body::wrap_stream(stream))
282                .send()
283                .await
284                .map_err(|_| unreachable_store())?;
285
286            offset += this;
287            let finished = offset == length;
288            match response.status().as_u16() {
289                RESUME_INCOMPLETE if !finished => {}
290                _ if finished => {
291                    expect_success(response, "finish a resumable upload").await?;
292                }
293                _ => {
294                    expect_success(response, "write a chunk").await?;
295                    return Err(malformed(
296                        "the object store finished an upload before it had every chunk",
297                    ));
298                }
299            }
300        }
301
302        Ok(())
303    }
304
305    pub(crate) async fn put_if_absent(&self, key: &str, body: Vec<u8>) -> Result<bool, Error> {
306        let length = body.len() as u64;
307        let response = self
308            .insert(key, reqwest::Body::from(body), length, true)
309            .await?;
310
311        if response.status() == StatusCode::PRECONDITION_FAILED {
312            return Ok(false);
313        }
314
315        expect_success(response, "write").await?;
316
317        Ok(true)
318    }
319
320    pub(crate) async fn get_bytes(&self, key: &str) -> Result<Option<Vec<u8>>, Error> {
321        let url = self.url(false, Some(key), &[("alt", "media")]);
322        let response = read_retrying(self.authorized(Method::GET, url).await?).await?;
323
324        if response.status() == StatusCode::NOT_FOUND {
325            return Ok(None);
326        }
327
328        expect_success(response, "read")
329            .await?
330            .bytes()
331            .await
332            .map(|bytes| Some(bytes.to_vec()))
333            .map_err(|_| unreachable_store())
334    }
335
336    pub(crate) async fn delete(&self, key: &str) -> Result<bool, Error> {
337        let url = self.url(false, Some(key), &[]);
338        let response = self
339            .authorized(Method::DELETE, url)
340            .await?
341            .send()
342            .await
343            .map_err(|_| unreachable_store())?;
344
345        if response.status() == StatusCode::NOT_FOUND {
346            return Ok(false);
347        }
348
349        expect_success(response, "delete").await?;
350
351        Ok(true)
352    }
353
354    pub(crate) async fn entries(&self, prefix: &str) -> Result<Vec<Entry>, Error> {
355        let mut out = Vec::new();
356        let mut token: Option<String> = None;
357
358        loop {
359            let mut query = vec![
360                ("prefix", prefix),
361                ("fields", "items(name,size,updated),nextPageToken"),
362            ];
363            if let Some(token) = &token {
364                query.push(("pageToken", token));
365            }
366
367            let url = self.url(false, None, &query);
368            let response = read_retrying(self.authorized(Method::GET, url).await?).await?;
369            let page: Page = expect_success(response, "list")
370                .await?
371                .json()
372                .await
373                .map_err(|_| {
374                    malformed("the object store sent a listing this server could not read")
375                })?;
376
377            for object in page.items {
378                out.push(Entry {
379                    size: object
380                        .size
381                        .parse()
382                        .map_err(|_| malformed("the object store listed an object with no size"))?,
383                    modified: object.updated.as_deref().and_then(|updated| {
384                        time::OffsetDateTime::parse(
385                            updated,
386                            &time::format_description::well_known::Rfc3339,
387                        )
388                        .ok()
389                    }),
390                    key: object.name,
391                });
392            }
393
394            match page.next_page_token.filter(|token| !token.is_empty()) {
395                Some(next) => token = Some(next),
396                None => break,
397            }
398        }
399
400        Ok(out)
401    }
402}
403
404#[cfg(test)]
405mod tests;