| 1 | //! 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 | |
| 6 | use worker::{Bucket, Env, Range, Result, UploadedPart}; |
| 7 | |
| 8 | use crate::sigv4::{Credentials, amz_date}; |
| 9 | use crate::{BlobStore, Config, Got, Part, Wanted, var}; |
| 10 | |
| 11 | pub struct R2Store { |
| 12 | bucket: Bucket, |
| 13 | signer: Option<(Credentials, String, String)>, |
| 14 | } |
| 15 | |
| 16 | impl 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 | |
| 33 | impl 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 | } |