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