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;