g1t/crates/blobstore/src/r2.rs

90 lines3,486 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Object storage is a shared port: packages' R2 and S3 adapters move to crates/blobstore1//! The R2 adapter: the service's bucket binding for everything, and R2's S3
2//! endpoint only to sign download URLs, when the service names signing
3//! variables (`Config::r2_signer`) and they are set. Without them, large
4//! objects stream through the Worker like small ones.
5
6use worker::{Bucket, Env, Range, Result, UploadedPart};
7
8use crate::sigv4::{Credentials, amz_date};
9use crate::{BlobStore, Config, Got, Part, Wanted, var};
10
11pub struct R2Store {
12 bucket: Bucket,
13 signer: Option<(Credentials, String, String)>,
14}
15
16impl R2Store {
17 pub fn from_env(env: &Env, config: &Config) -> Result<R2Store> {
18 let [key, secret, account, bucket] = config
19 .r2_signer
20 .map(|names| names.map(|name| var(env, name)))
21 .unwrap_or_default();
22 let signer = (!key.is_empty() && !secret.is_empty() && !account.is_empty() && !bucket.is_empty()).then(|| {
23 (
24 Credentials { access_key_id: key, secret_access_key: secret, region: "auto".to_owned() },
25 format!("{account}.r2.cloudflarestorage.com"),
26 bucket,
27 )
28 });
29 Ok(R2Store { bucket: env.bucket(config.binding)?, signer })
30 }
31}
32
33impl BlobStore for R2Store {
34 async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
35 self.bucket.put(key, bytes).execute().await?;
36 Ok(())
37 }
38
39 async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>> {
40 let mut get = self.bucket.get(key);
41 if let Some(range) = range {
42 get = get.range(Range::OffsetWithLength { offset: range.offset, length: range.length });
43 }
44 let Some(object) = get.execute().await? else {
45 return Ok(None);
46 };
47 let size = object.size();
48 let Some(body) = object.body() else {
49 return Ok(None);
50 };
51 Ok(Some(Got { size, body: body.response_body()? }))
52 }
53
54 async fn head(&self, key: &str) -> Result<Option<u64>> {
55 Ok(self.bucket.head(key).await?.map(|object| object.size()))
56 }
57
58 async fn delete(&self, key: &str) -> Result<()> {
59 self.bucket.delete(key).await
60 }
61
62 async fn create_multipart(&self, key: &str) -> Result<String> {
63 let upload = self.bucket.create_multipart_upload(key).execute().await?;
64 Ok(upload.upload_id().await)
65 }
66
67 async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part> {
68 let upload = self.bucket.resume_multipart_upload(key, upload_id)?;
69 let part = upload.upload_part(number, bytes).await?;
70 Ok(Part { number: part.part_number(), etag: part.etag() })
71 }
72
73 async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()> {
74 let upload = self.bucket.resume_multipart_upload(key, upload_id)?;
75 upload
76 .complete(parts.iter().map(|part| UploadedPart::new(part.number, part.etag.clone())))
77 .await?;
78 Ok(())
79 }
80
81 async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()> {
82 self.bucket.resume_multipart_upload(key, upload_id)?.abort().await
83 }
84
85 fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String> {
86 let (credentials, host, bucket) = self.signer.as_ref()?;
87 let path = format!("/{bucket}/{key}");
88 Some(credentials.presign_get(&format!("https://{host}"), host, &path, &amz_date(now_ms), expires))
89 }
90}