g1t/services/packages/src/store/r2.rs

93 lines3,521 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.

Packages, with a container registry on g1t.sh; workspaces deleted whole and kept 30 days; Members for every member1//! 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
6use worker::{Bucket, Env, Range, Result, UploadedPart};
7
8use super::{BlobStore, Got, Part, var};
9use crate::range::Wanted;
10use crate::sigv4::{Credentials, amz_date};
11
12pub struct R2Store {
13 bucket: Bucket,
14 signer: Option<(Credentials, String, String)>,
15}
16
17impl 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
36impl 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}