Commit

Object storage is a shared port: packages' R2 and S3 adapters move to crates/blobstore

Each service names its own bucket binding and variables (Config), so the repos service can keep backups through the same adapters.

syntaqxcommitted Parent3348046Browse files
15 files+803−7160/15 viewed
+14−0
941941 ]
942942
943943 [[package]]
944+name = "g1t-blobstore"
945+version = "0.1.0"
946+dependencies = [
947+ "g1t-contracts",
948+ "g1t-kit",
949+ "hex",
950+ "hmac 0.12.1",
951+ "serde",
952+ "sha2 0.10.9",
953+ "worker",
954+]
955+
956+[[package]]
944957 name = "g1t-contracts"
945958 version = "0.1.0"
946959 dependencies = [
10091022 dependencies = [
10101023 "base64 0.22.1",
10111024 "futures-util",
1025+ "g1t-blobstore",
10121026 "g1t-contracts",
10131027 "g1t-kit",
10141028 "hex",
+1−0
99
1010 [workspace.dependencies]
1111 g1t-actions = { path = "crates/actions" }
12+g1t-blobstore = { path = "crates/blobstore" }
1213 g1t-contracts = { path = "crates/contracts" }
1314 g1t-kit = { path = "crates/kit" }
1415 g1t-scan = { path = "crates/scan" }
+15−0
1+[package]
2+name = "g1t-blobstore"
3+version = "0.1.0"
4+edition.workspace = true
5+license.workspace = true
6+description = "Object storage behind one port: R2 on Cloudflare, any S3-compatible store (MinIO) when self-hosted."
7+
8+[dependencies]
9+g1t-contracts.workspace = true
10+g1t-kit.workspace = true
11+serde.workspace = true
12+worker.workspace = true
13+hex = "0.4"
14+hmac = "0.12"
15+sha2 = { version = "0.10", features = ["compress"] }
+192−0
1+//! Object storage behind one port, `BlobStore`: R2 on Cloudflare, and any
2+//! S3-compatible store (MinIO in the self-host compose file) elsewhere.
3+//!
4+//! Every service that keeps objects names its own [`Config`]: the variable
5+//! that chooses the store (`r2`, the default, or `s3`), the R2 bucket
6+//! binding, and the variable naming its S3 bucket. The S3 endpoint and its
7+//! credentials (S3_ENDPOINT, S3_REGION, S3_ACCESS_KEY_ID,
8+//! S3_SECRET_ACCESS_KEY) are the installation's, shared by every service.
9+//!
10+//! Large objects go up as multipart parts of one size, as R2 requires
11+//! (every part but the last the same size).
12+
13+mod r2;
14+mod s3;
15+pub mod sigv4;
16+
17+use serde::{Deserialize, Serialize};
18+use worker::{Env, Response, ResponseBody, Result};
19+
20+pub use r2::R2Store;
21+pub use s3::S3Store;
22+
23+/// A part of an object a read asks for: `length` bytes from `offset`.
24+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
25+pub struct Wanted {
26+ pub offset: u64,
27+ pub length: u64,
28+}
29+
30+impl Wanted {
31+ /// `bytes <first>-<last>/<size>`.
32+ pub fn content_range(&self, size: u64) -> String {
33+ format!("bytes {}-{}/{size}", self.offset, self.offset + self.length - 1)
34+ }
35+}
36+
37+/// One part of a multipart upload, as completing it needs.
38+#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
39+pub struct Part {
40+ pub number: u16,
41+ pub etag: String,
42+}
43+
44+/// An object read back.
45+pub struct Got {
46+ /// The whole object's size, whatever range was read.
47+ pub size: u64,
48+ pub body: ResponseBody,
49+}
50+
51+impl Got {
52+ pub async fn bytes(self) -> Result<Vec<u8>> {
53+ match self.body {
54+ ResponseBody::Empty => Ok(Vec::new()),
55+ ResponseBody::Body(bytes) => Ok(bytes),
56+ stream => Response::from_body(stream)?.bytes().await,
57+ }
58+ }
59+}
60+
61+/// What a service needs of storage.
62+#[allow(async_fn_in_trait)]
63+pub trait BlobStore {
64+ async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()>;
65+ /// The object, or the part of it `range` asks for.
66+ async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>>;
67+ /// The object's size, if it is there.
68+ async fn head(&self, key: &str) -> Result<Option<u64>>;
69+ async fn delete(&self, key: &str) -> Result<()>;
70+ /// Starts a multipart upload to `key`, and says its id.
71+ async fn create_multipart(&self, key: &str) -> Result<String>;
72+ async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part>;
73+ async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()>;
74+ async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()>;
75+ /// A URL that downloads the object for `expires` seconds without
76+ /// passing through this Worker, when the store can sign one.
77+ fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String>;
78+
79+ /// The whole object, read into memory: for small ones only.
80+ async fn read(&self, key: &str) -> Result<Option<Vec<u8>>> {
81+ match self.get(key, None).await? {
82+ Some(got) => Ok(Some(got.bytes().await?)),
83+ None => Ok(None),
84+ }
85+ }
86+}
87+
88+/// Where one service's objects are, by the names of its bindings and
89+/// variables.
90+#[derive(Clone, Copy, Debug)]
91+pub struct Config {
92+ /// The variable that chooses the store: `r2` (or unset) or `s3`.
93+ pub kind: &'static str,
94+ /// The R2 bucket binding.
95+ pub binding: &'static str,
96+ /// The variables that let R2's S3 endpoint sign download URLs: access
97+ /// key id, secret, account id and bucket name. None: never signed.
98+ pub r2_signer: Option<[&'static str; 4]>,
99+ /// The variable naming the S3 bucket.
100+ pub s3_bucket: &'static str,
101+ /// The variable naming where clients reach the S3 store, for signed
102+ /// downloads. None: never signed.
103+ pub s3_public_endpoint: Option<&'static str>,
104+}
105+
106+/// The store a service is configured with.
107+pub enum Store {
108+ R2(R2Store),
109+ S3(S3Store),
110+}
111+
112+impl Store {
113+ pub fn from_env(env: &Env, config: &Config) -> Result<Store> {
114+ if var(env, config.kind) == "s3" {
115+ return Ok(Store::S3(S3Store::from_env(env, config)?));
116+ }
117+ Ok(Store::R2(R2Store::from_env(env, config)?))
118+ }
119+}
120+
121+impl BlobStore for Store {
122+ async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
123+ match self {
124+ Store::R2(s) => s.put(key, bytes).await,
125+ Store::S3(s) => s.put(key, bytes).await,
126+ }
127+ }
128+
129+ async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>> {
130+ match self {
131+ Store::R2(s) => s.get(key, range).await,
132+ Store::S3(s) => s.get(key, range).await,
133+ }
134+ }
135+
136+ async fn head(&self, key: &str) -> Result<Option<u64>> {
137+ match self {
138+ Store::R2(s) => s.head(key).await,
139+ Store::S3(s) => s.head(key).await,
140+ }
141+ }
142+
143+ async fn delete(&self, key: &str) -> Result<()> {
144+ match self {
145+ Store::R2(s) => s.delete(key).await,
146+ Store::S3(s) => s.delete(key).await,
147+ }
148+ }
149+
150+ async fn create_multipart(&self, key: &str) -> Result<String> {
151+ match self {
152+ Store::R2(s) => s.create_multipart(key).await,
153+ Store::S3(s) => s.create_multipart(key).await,
154+ }
155+ }
156+
157+ async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part> {
158+ match self {
159+ Store::R2(s) => s.upload_part(key, upload_id, number, bytes).await,
160+ Store::S3(s) => s.upload_part(key, upload_id, number, bytes).await,
161+ }
162+ }
163+
164+ async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()> {
165+ match self {
166+ Store::R2(s) => s.complete_multipart(key, upload_id, parts).await,
167+ Store::S3(s) => s.complete_multipart(key, upload_id, parts).await,
168+ }
169+ }
170+
171+ async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()> {
172+ match self {
173+ Store::R2(s) => s.abort_multipart(key, upload_id).await,
174+ Store::S3(s) => s.abort_multipart(key, upload_id).await,
175+ }
176+ }
177+
178+ fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String> {
179+ match self {
180+ Store::R2(s) => s.presign_get(key, expires, now_ms),
181+ Store::S3(s) => s.presign_get(key, expires, now_ms),
182+ }
183+ }
184+}
185+
186+/// A configuration variable or secret, or the empty string.
187+pub fn var(env: &Env, name: &str) -> String {
188+ env.var(name)
189+ .map(|v| v.to_string())
190+ .or_else(|_| env.secret(name).map(|v| v.to_string()))
191+ .unwrap_or_default()
192+}
+90−0
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+}
+255−0
1+//! The S3 adapter, for self-hosted installations: any S3-compatible store
2+//! (MinIO, Ceph, Garage, AWS) over fetch, signed with SigV4, path-style.
3+//! S3_ENDPOINT, S3_ACCESS_KEY_ID, S3_SECRET_ACCESS_KEY and S3_REGION say
4+//! where, and the service's own variable (`Config::s3_bucket`) which
5+//! bucket. Its public endpoint variable, when it names one and that is
6+//! set, is the address clients reach the store at, and large downloads are
7+//! then sent there with a signed URL instead of through the Worker.
8+
9+use worker::wasm_bindgen::JsValue;
10+use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Response, Result, Url};
11+
12+use crate::sigv4::{Credentials, UNSIGNED, amz_date};
13+use crate::{BlobStore, Config, Got, Part, Wanted, var};
14+
15+pub struct S3Store {
16+ /// `http://minio:9000`, without a trailing slash.
17+ endpoint: String,
18+ /// Where clients reach the same store, for signed URLs.
19+ public_endpoint: Option<String>,
20+ bucket: String,
21+ credentials: Credentials,
22+}
23+
24+fn failed(what: &str, status: u16, body: &str) -> worker::Error {
25+ let said: String = body.chars().take(300).collect();
26+ worker::Error::RustError(format!("storage {what} failed with status {status}: {said}"))
27+}
28+
29+/// The text of the first `<tag>` in an XML answer.
30+fn xml_value<'a>(xml: &'a str, tag: &str) -> Option<&'a str> {
31+ let open = format!("<{tag}>");
32+ let start = xml.find(&open)? + open.len();
33+ let end = xml[start..].find(&format!("</{tag}>"))? + start;
34+ Some(&xml[start..end])
35+}
36+
37+fn host_of(endpoint: &str) -> String {
38+ endpoint
39+ .split_once("://")
40+ .map_or(endpoint, |(_, rest)| rest)
41+ .split('/')
42+ .next()
43+ .unwrap_or_default()
44+ .to_owned()
45+}
46+
47+impl S3Store {
48+ pub fn from_env(env: &Env, config: &Config) -> Result<S3Store> {
49+ let endpoint = var(env, "S3_ENDPOINT").trim_end_matches('/').to_owned();
50+ let bucket = var(env, config.s3_bucket);
51+ if endpoint.is_empty() || bucket.is_empty() {
52+ return Err(worker::Error::RustError(format!(
53+ "{} is s3, but S3_ENDPOINT or {} is not set",
54+ config.kind, config.s3_bucket
55+ )));
56+ }
57+ let region = var(env, "S3_REGION");
58+ let public = config
59+ .s3_public_endpoint
60+ .map(|name| var(env, name).trim_end_matches('/').to_owned())
61+ .unwrap_or_default();
62+ Ok(S3Store {
63+ endpoint,
64+ public_endpoint: (!public.is_empty()).then_some(public),
65+ bucket,
66+ credentials: Credentials {
67+ access_key_id: var(env, "S3_ACCESS_KEY_ID"),
68+ secret_access_key: var(env, "S3_SECRET_ACCESS_KEY"),
69+ region: if region.is_empty() { "us-east-1".to_owned() } else { region },
70+ },
71+ })
72+ }
73+
74+ fn path(&self, key: &str) -> String {
75+ format!("/{}/{key}", self.bucket)
76+ }
77+
78+ /// Sends one signed request, and answers with the response whatever
79+ /// its status.
80+ async fn send(
81+ &self,
82+ method: Method,
83+ key: &str,
84+ query: &[(String, String)],
85+ extra: &[(&str, String)],
86+ body: Option<Vec<u8>>,
87+ ) -> Result<Response> {
88+ let path = self.path(key);
89+ let date = amz_date(g1t_kit::now_ms());
90+ let mut signed = vec![
91+ ("host".to_owned(), host_of(&self.endpoint)),
92+ ("x-amz-content-sha256".to_owned(), UNSIGNED.to_owned()),
93+ ("x-amz-date".to_owned(), date),
94+ ];
95+ for (name, value) in extra {
96+ signed.push(((*name).to_owned(), value.clone()));
97+ }
98+ let authorization = self
99+ .credentials
100+ .authorization(method.as_ref(), &path, query, &signed, UNSIGNED);
101+ let headers = Headers::new();
102+ for (name, value) in &signed {
103+ if name != "host" {
104+ headers.set(name, value)?;
105+ }
106+ }
107+ headers.set("authorization", &authorization)?;
108+ let mut url = Url::parse(&format!("{}{}", self.endpoint, crate::sigv4::uri_encode(&path, true)))?;
109+ if !query.is_empty() {
110+ let text: Vec<String> = query
111+ .iter()
112+ .map(|(k, v)| {
113+ let (k, v) = (crate::sigv4::uri_encode(k, false), crate::sigv4::uri_encode(v, false));
114+ if v.is_empty() { format!("{k}=") } else { format!("{k}={v}") }
115+ })
116+ .collect();
117+ url.set_query(Some(&text.join("&")));
118+ }
119+ let mut init = RequestInit::new();
120+ init.with_method(method).with_headers(headers);
121+ if let Some(body) = body {
122+ init.with_body(Some(JsValue::from(worker::js_sys::Uint8Array::from(body.as_slice()))));
123+ }
124+ Fetch::Request(Request::new_with_init(url.as_str(), &init)?).send().await
125+ }
126+
127+ async fn ok(&self, what: &str, mut response: Response) -> Result<Response> {
128+ let status = response.status_code();
129+ if (200..300).contains(&status) {
130+ return Ok(response);
131+ }
132+ let body = response.text().await.unwrap_or_default();
133+ Err(failed(what, status, &body))
134+ }
135+}
136+
137+fn query(pairs: &[(&str, &str)]) -> Vec<(String, String)> {
138+ pairs.iter().map(|(k, v)| ((*k).to_owned(), (*v).to_owned())).collect()
139+}
140+
141+impl BlobStore for S3Store {
142+ async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
143+ let length = bytes.len().to_string();
144+ let response = self
145+ .send(Method::Put, key, &[], &[("content-length", length)], Some(bytes))
146+ .await?;
147+ self.ok("put", response).await.map(|_| ())
148+ }
149+
150+ async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>> {
151+ let extra: Vec<(&str, String)> = range
152+ .map(|r| ("range", format!("bytes={}-{}", r.offset, r.offset + r.length - 1)))
153+ .into_iter()
154+ .collect();
155+ let response = self.send(Method::Get, key, &[], &extra, None).await?;
156+ if response.status_code() == 404 {
157+ return Ok(None);
158+ }
159+ let response = self.ok("get", response).await?;
160+ let size = match response.headers().get("content-range")? {
161+ // `bytes 0-9/100`: the whole object's size is after the slash.
162+ Some(range) => range.rsplit('/').next().and_then(|n| n.parse().ok()).unwrap_or(0),
163+ None => response.headers().get("content-length")?.and_then(|n| n.parse().ok()).unwrap_or(0),
164+ };
165+ let (_, body) = response.into_parts();
166+ Ok(Some(Got { size, body }))
167+ }
168+
169+ async fn head(&self, key: &str) -> Result<Option<u64>> {
170+ let response = self.send(Method::Head, key, &[], &[], None).await?;
171+ if response.status_code() == 404 {
172+ return Ok(None);
173+ }
174+ let response = self.ok("head", response).await?;
175+ Ok(response.headers().get("content-length")?.and_then(|n| n.parse().ok()))
176+ }
177+
178+ async fn delete(&self, key: &str) -> Result<()> {
179+ let response = self.send(Method::Delete, key, &[], &[], None).await?;
180+ if response.status_code() == 404 {
181+ return Ok(());
182+ }
183+ self.ok("delete", response).await.map(|_| ())
184+ }
185+
186+ async fn create_multipart(&self, key: &str) -> Result<String> {
187+ let response = self.send(Method::Post, key, &query(&[("uploads", "")]), &[], None).await?;
188+ let text = self.ok("create multipart", response).await?.text().await?;
189+ xml_value(&text, "UploadId")
190+ .map(str::to_owned)
191+ .ok_or_else(|| failed("create multipart", 200, &text))
192+ }
193+
194+ async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part> {
195+ let number_text = number.to_string();
196+ let length = bytes.len().to_string();
197+ let response = self
198+ .send(
199+ Method::Put,
200+ key,
201+ &query(&[("partNumber", &number_text), ("uploadId", upload_id)]),
202+ &[("content-length", length)],
203+ Some(bytes),
204+ )
205+ .await?;
206+ let response = self.ok("upload part", response).await?;
207+ let etag = response.headers().get("etag")?.unwrap_or_default();
208+ Ok(Part { number, etag })
209+ }
210+
211+ async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()> {
212+ let mut xml = String::from("<CompleteMultipartUpload>");
213+ for part in parts {
214+ xml.push_str(&format!("<Part><PartNumber>{}</PartNumber><ETag>{}</ETag></Part>", part.number, part.etag));
215+ }
216+ xml.push_str("</CompleteMultipartUpload>");
217+ let length = xml.len().to_string();
218+ let response = self
219+ .send(Method::Post, key, &query(&[("uploadId", upload_id)]), &[("content-length", length)], Some(xml.into_bytes()))
220+ .await?;
221+ // S3 may answer 200 and still have failed, saying so in the body.
222+ let text = self.ok("complete multipart", response).await?.text().await?;
223+ if text.contains("<Error>") {
224+ return Err(failed("complete multipart", 200, &text));
225+ }
226+ Ok(())
227+ }
228+
229+ async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()> {
230+ let response = self.send(Method::Delete, key, &query(&[("uploadId", upload_id)]), &[], None).await?;
231+ if response.status_code() == 404 {
232+ return Ok(());
233+ }
234+ self.ok("abort multipart", response).await.map(|_| ())
235+ }
236+
237+ fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String> {
238+ let base = self.public_endpoint.as_ref()?;
239+ Some(self.credentials.presign_get(base, &host_of(base), &self.path(key), &amz_date(now_ms), expires))
240+ }
241+}
242+
243+#[cfg(test)]
244+mod tests {
245+ use super::*;
246+
247+ #[test]
248+ fn answers_are_read_from_their_xml() {
249+ let xml = "<InitiateMultipartUploadResult><Bucket>b</Bucket><UploadId>abc-123</UploadId></InitiateMultipartUploadResult>";
250+ assert_eq!(xml_value(xml, "UploadId"), Some("abc-123"));
251+ assert_eq!(xml_value(xml, "Key"), None);
252+ assert_eq!(host_of("http://minio:9000"), "minio:9000");
253+ assert_eq!(host_of("https://s3.example.com/base"), "s3.example.com");
254+ }
255+}
+201−0
1+//! AWS Signature Version 4, for S3-compatible storage: signing a request's
2+//! headers, and signing a URL that lets its holder download one object for
3+//! a while. R2's S3 endpoint takes the same signatures, which is how large
4+//! downloads are sent straight to it.
5+
6+use hmac::{Hmac, Mac};
7+use sha2::{Digest as _, Sha256};
8+
9+type HmacSha256 = Hmac<Sha256>;
10+
11+/// The hash of a payload that is not signed: bodies stream as they are.
12+pub const UNSIGNED: &str = "UNSIGNED-PAYLOAD";
13+
14+/// Who signs, and for which region.
15+#[derive(Clone, Debug)]
16+pub struct Credentials {
17+ pub access_key_id: String,
18+ pub secret_access_key: String,
19+ pub region: String,
20+}
21+
22+/// `20130524T000000Z` from milliseconds since the epoch.
23+pub fn amz_date(now_ms: u64) -> String {
24+ let text = g1t_contracts::time::rfc3339(now_ms);
25+ let whole = text.split('.').next().unwrap_or(&text);
26+ format!("{}Z", whole.replace(['-', ':'], ""))
27+}
28+
29+/// Percent-encodes everything but the unreserved characters, and `/` too
30+/// unless `path`.
31+pub fn uri_encode(text: &str, path: bool) -> String {
32+ let mut out = String::with_capacity(text.len());
33+ for byte in text.bytes() {
34+ match byte {
35+ b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => out.push(byte as char),
36+ b'/' if path => out.push('/'),
37+ _ => out.push_str(&format!("%{byte:02X}")),
38+ }
39+ }
40+ out
41+}
42+
43+fn hmac(key: &[u8], data: &str) -> Vec<u8> {
44+ let mut mac = HmacSha256::new_from_slice(key).expect("HMAC takes a key of any length");
45+ mac.update(data.as_bytes());
46+ mac.finalize().into_bytes().to_vec()
47+}
48+
49+fn sha256_hex(data: &[u8]) -> String {
50+ hex::encode(Sha256::digest(data))
51+}
52+
53+/// The query string, sorted and encoded as signing needs it.
54+fn canonical_query(query: &[(String, String)]) -> String {
55+ let mut pairs: Vec<(String, String)> = query
56+ .iter()
57+ .map(|(key, value)| (uri_encode(key, false), uri_encode(value, false)))
58+ .collect();
59+ pairs.sort();
60+ pairs
61+ .iter()
62+ .map(|(key, value)| format!("{key}={value}"))
63+ .collect::<Vec<_>>()
64+ .join("&")
65+}
66+
67+impl Credentials {
68+ fn scope(&self, date: &str) -> String {
69+ format!("{}/{}/s3/aws4_request", &date[..8], self.region)
70+ }
71+
72+ fn signature(&self, date: &str, canonical_request: &str) -> String {
73+ let to_sign = format!(
74+ "AWS4-HMAC-SHA256\n{date}\n{}\n{}",
75+ self.scope(date),
76+ sha256_hex(canonical_request.as_bytes())
77+ );
78+ let key = hmac(format!("AWS4{}", self.secret_access_key).as_bytes(), &date[..8]);
79+ let key = hmac(&key, &self.region);
80+ let key = hmac(&key, "s3");
81+ let key = hmac(&key, "aws4_request");
82+ hex::encode(hmac(&key, &to_sign))
83+ }
84+
85+ /// The `Authorization` header for a request. `headers` must include
86+ /// `host`, `x-amz-date` and `x-amz-content-sha256`, lowercase; every
87+ /// one given is signed.
88+ pub fn authorization(
89+ &self,
90+ method: &str,
91+ path: &str,
92+ query: &[(String, String)],
93+ headers: &[(String, String)],
94+ payload_hash: &str,
95+ ) -> String {
96+ let mut headers: Vec<(String, String)> = headers
97+ .iter()
98+ .map(|(name, value)| (name.to_ascii_lowercase(), value.trim().to_owned()))
99+ .collect();
100+ headers.sort();
101+ let date = headers
102+ .iter()
103+ .find(|(name, _)| name == "x-amz-date")
104+ .map(|(_, value)| value.clone())
105+ .unwrap_or_default();
106+ let signed: Vec<&str> = headers.iter().map(|(name, _)| name.as_str()).collect();
107+ let signed = signed.join(";");
108+ let canonical_headers: String = headers.iter().map(|(name, value)| format!("{name}:{value}\n")).collect();
109+ let canonical = format!(
110+ "{method}\n{}\n{}\n{canonical_headers}\n{signed}\n{payload_hash}",
111+ uri_encode(path, true),
112+ canonical_query(query)
113+ );
114+ format!(
115+ "AWS4-HMAC-SHA256 Credential={}/{}, SignedHeaders={signed}, Signature={}",
116+ self.access_key_id,
117+ self.scope(&date),
118+ self.signature(&date, &canonical)
119+ )
120+ }
121+
122+ /// A URL that lets anyone `GET` the object at `path` on `host` for
123+ /// `expires` seconds from `date`. `base` is the scheme and host the
124+ /// URL starts with.
125+ pub fn presign_get(&self, base: &str, host: &str, path: &str, date: &str, expires: u32) -> String {
126+ let mut query = vec![
127+ ("X-Amz-Algorithm".to_owned(), "AWS4-HMAC-SHA256".to_owned()),
128+ ("X-Amz-Credential".to_owned(), format!("{}/{}", self.access_key_id, self.scope(date))),
129+ ("X-Amz-Date".to_owned(), date.to_owned()),
130+ ("X-Amz-Expires".to_owned(), expires.to_string()),
131+ ("X-Amz-SignedHeaders".to_owned(), "host".to_owned()),
132+ ];
133+ let canonical = format!(
134+ "GET\n{}\n{}\nhost:{host}\n\nhost\n{UNSIGNED}",
135+ uri_encode(path, true),
136+ canonical_query(&query)
137+ );
138+ query.push(("X-Amz-Signature".to_owned(), self.signature(date, &canonical)));
139+ format!("{base}{}?{}", uri_encode(path, true), canonical_query(&query))
140+ }
141+}
142+
143+#[cfg(test)]
144+mod tests {
145+ use super::*;
146+
147+ fn example() -> Credentials {
148+ Credentials {
149+ access_key_id: "AKIAIOSFODNN7EXAMPLE".into(),
150+ secret_access_key: "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY".into(),
151+ region: "us-east-1".into(),
152+ }
153+ }
154+
155+ /// AWS's own example of a presigned URL (Authenticating Requests:
156+ /// Using Query Parameters).
157+ #[test]
158+ fn a_presigned_url_matches_the_aws_example() {
159+ let url = example().presign_get(
160+ "https://examplebucket.s3.amazonaws.com",
161+ "examplebucket.s3.amazonaws.com",
162+ "/test.txt",
163+ "20130524T000000Z",
164+ 86400,
165+ );
166+ assert!(url.starts_with("https://examplebucket.s3.amazonaws.com/test.txt?X-Amz-Algorithm=AWS4-HMAC-SHA256"));
167+ assert!(url.contains("X-Amz-Credential=AKIAIOSFODNN7EXAMPLE%2F20130524%2Fus-east-1%2Fs3%2Faws4_request"));
168+ assert!(url.contains("&X-Amz-Signature=aeeed9bbccd4d02ee5c0109b86d86835f995330da4c265957d157751f604d404&"), "{url}");
169+ }
170+
171+ /// AWS's own example of a signed GET with a range (Authenticating
172+ /// Requests: Using the Authorization Header).
173+ #[test]
174+ fn a_signed_request_matches_the_aws_example() {
175+ let empty = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855";
176+ let headers = [
177+ ("Host", "examplebucket.s3.amazonaws.com"),
178+ ("Range", "bytes=0-9"),
179+ ("x-amz-content-sha256", empty),
180+ ("x-amz-date", "20130524T000000Z"),
181+ ]
182+ .map(|(name, value)| (name.to_owned(), value.to_owned()));
183+ let authorization = example().authorization("GET", "/test.txt", &[], &headers, empty);
184+ assert_eq!(
185+ authorization,
186+ "AWS4-HMAC-SHA256 Credential=AKIAIOSFODNN7EXAMPLE/20130524/us-east-1/s3/aws4_request, \
187+ SignedHeaders=host;range;x-amz-content-sha256;x-amz-date, \
188+ Signature=f0e8bdb87c964420e857bd35b5d6ed310bd44f0170aba48dd91039c6036bdb41"
189+ );
190+ }
191+
192+ #[test]
193+ fn dates_and_encoding() {
194+ assert_eq!(amz_date(1_369_353_600_000), "20130524T000000Z");
195+ assert_eq!(uri_encode("a b/c+d~", true), "a%20b/c%2Bd~");
196+ assert_eq!(uri_encode("a/b", false), "a%2Fb");
197+ let query = [("uploadId".to_owned(), "x y".to_owned()), ("partNumber".to_owned(), "2".to_owned())];
198+ assert_eq!(canonical_query(&query), "partNumber=2&uploadId=x%20y");
199+ assert_eq!(canonical_query(&[("uploads".to_owned(), String::new())]), "uploads=");
200+ }
201+}
+1−0
1111 [dependencies]
1212 g1t-contracts.workspace = true
1313 g1t-kit.workspace = true
14+g1t-blobstore.workspace = true
1415 serde.workspace = true
1516 serde_json.workspace = true
1617 worker.workspace = true
+1−2
2222 mod oci;
2323 mod quota;
2424 mod range;
25−mod sigv4;
2625 mod store;
2726 mod token;
2827 mod upload;
166165 let host = store::var(env, "REGISTRY_HOST");
167166 Ok(Packages {
168167 db: Db { db: env.d1("DB")? },
169− store: Store::from_env(env)?,
168+ store: store::from_env(env)?,
170169 identity: env.service("IDENTITY")?,
171170 repos: env.service("REPOS")?,
172171 events: env.service("EVENTS")?,
+1−12
2727 }
2828
2929 /// A part of a blob a download asks for, resolved against its size.
30−#[derive(Clone, Copy, Debug, PartialEq, Eq)]
31−pub struct Wanted {
32− pub offset: u64,
33− pub length: u64,
34−}
35−
36−impl Wanted {
37− /// `bytes <first>-<last>/<size>`.
38− pub fn content_range(&self, size: u64) -> String {
39− format!("bytes {}-{}/{size}", self.offset, self.offset + self.length - 1)
40− }
41−}
30+pub use g1t_blobstore::Wanted;
4231
4332 /// What a download's `Range` header asks for, against a blob of `size`
4433 /// bytes. `Ok(None)`: the whole blob (no header, or one this does not
+0−201
1−//! AWS Signature Version 4, for S3-compatible storage: signing a request's
2−//! headers, and signing a URL that lets its holder download one object for
3−//! a while. R2's S3 endpoint takes the same signatures, which is how large
4−//! downloads are sent straight to it.
5−
6−use hmac::{Hmac, Mac};
7−use sha2::{Digest as _, Sha256};
8−
9−type HmacSha256 = Hmac<Sha256>;
10−
11−/// The hash of a payload that is not signed: bodies stream as they are.
12−pub const UNSIGNED: &str = "UNSIGNED-PAYLOAD";
13−
14−/// Who signs, and for which region.
15−#[derive(Clone, Debug)]
16−pub struct Credentials {
17− pub access_key_id: String,
18− pub secret_access_key: String,
19− pub region: String,
20−}
21−
22−/// `20130524T000000Z` from milliseconds since the epoch.
23−pub fn amz_date(now_ms: u64) -> String {
24− let text = g1t_contracts::time::rfc3339(now_ms);
25− let whole = text.split('.').next().unwrap_or(&text);
26− format!("{}Z", whole.replace(['-', ':'], ""))
27−}
28−
29−/// Percent-encodes everything but the unreserved characters, and `/` too
30−/// unless `path`.
31−pub fn uri_encode(text: &str, path: bool) -> String {
32− let mut out = String::with_capacity(text.len());
33− for byte in text.bytes() {
34− match byte {
35− b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => out.push(byte as char),
36− b'/' if path => out.push('/'),
37− _ => out.push_str(&format!("%{byte:02X}")),
38− }
39− }
40− out
41−}
42−
43−fn hmac(key: &[u8], data: &str) -> Vec<u8> {
44− let mut mac = HmacSha256::new_from_slice(key).expect("HMAC takes a key of any length");
45− mac.update(data.as_bytes());
46− mac.finalize().into_bytes().to_vec()
47−}
48−
49−fn sha256_hex(data: &[u8]) -> String {
50− hex::encode(Sha256::digest(data))
51−}
52−
53−/// The query string, sorted and encoded as signing needs it.
54−fn canonical_query(query: &[(String, String)]) -> String {
55− let mut pairs: Vec<(String, String)> = query
56− .iter()
57− .map(|(key, value)| (uri_encode(key, false), uri_encode(value, false)))
58− .collect();
59− pairs.sort();
60− pairs
61− .iter()
62− .map(|(key, value)| format!("{key}={value}"))
63− .collect::<Vec<_>>()
64− .join("&")
65−}
66−
67−impl Credentials {
68− fn scope(&self, date: &str) -> String {
69− format!("{}/{}/s3/aws4_request", &date[..8], self.region)
70− }
71−
72− fn signature(&self, date: &str, canonical_request: &str) -> String {
73− let to_sign = format!(
74− "AWS4-HMAC-SHA256\n{date}\n{}\n{}",
75− self.scope(date),
76− sha256_hex(canonical_request.as_bytes())
77− );
78− let key = hmac(format!("AWS4{}", self.secret_access_key).as_bytes(), &date[..8]);
79− let key = hmac(&key, &self.region);
80− let key = hmac(&key, "s3");
81− let key = hmac(&key, "aws4_request");
82− hex::encode(hmac(&key, &to_sign))
83− }
84−
85− /// The `Authorization` header for a request. `headers` must include
86− /// `host`, `x-amz-date` and `x-amz-content-sha256`, lowercase; every
87− /// one given is signed.
88− pub fn authorization(
89− &self,
90− method: &str,
91− path: &str,
92− query: &[(String, String)],
93− headers: &[(String, String)],
94− payload_hash: &str,
95− ) -> String {
96− let mut headers: Vec<(String, String)> = headers
97− .iter()
98− .map(|(name, value)| (name.to_ascii_lowercase(), value.trim().to_owned()))
99− .collect();
100− headers.sort();
101− let date = headers
102− .iter()
103− .find(|(name, _)| name == "x-amz-date")
104− .map(|(_, value)| value.clone())
105− .unwrap_or_default();
106− let signed: Vec<&str> = headers.iter().map(|(name, _)| name.as_str()).collect();
107− let signed = signed.join(";");
108− let canonical_headers: String = headers.iter().map(|(name, value)| format!("{name}:{value}\n")).collect();
109− let canonical = format!(
110− "{method}\n{}\n{}\n{canonical_headers}\n{signed}\n{payload_hash}",
111− uri_encode(path, true),
112− canonical_query(query)
113− );
114− format!(
115− "AWS4-HMAC-SHA256 Credential={}/{}, SignedHeaders={signed}, Signature={}",
116− self.access_key_id,
117− self.scope(&date),
118− self.signature(&date, &canonical)
119− )
120− }
121−
122− /// A URL that lets anyone `GET` the object at `path` on `host` for
123− /// `expires` seconds from `date`. `base` is the scheme and host the
124− /// URL starts with.
125− pub fn presign_get(&self, base: &str, host: &str, path: &str, date: &str, expires: u32) -> String {
126− let mut query = vec![
127− ("X-Amz-Algorithm".to_owned(), "AWS4-HMAC-SHA256".to_owned()),
128− ("X-Amz-Credential".to_owned(), format!("{}/{}", self.access_key_id, self.scope(date))),
129− ("X-Amz-Date".to_owned(), date.to_owned()),
130− ("X-Amz-Expires".to_owned(), expires.to_string()),
131− ("X-Amz-SignedHeaders".to_owned(), "host".to_owned()),
132− ];
133− let canonical = format!(
134− "GET\n{}\n{}\nhost:{host}\n\nhost\n{UNSIGNED}",
135− uri_encode(path, true),
136− canonical_query(&query)
137− );
138− query.push(("X-Amz-Signature".to_owned(), self.signature(date, &canonical)));
139− format!("{base}{}?{}", uri_encode(path, true), canonical_query(&query))
140− }
141−}
142−
143−#[cfg(test)]
144−mod tests {
145− use super::*;
146−
147− fn example() -> Credentials {
148− Credentials {
149− access_key_id: "AKIAIOSFODNN7EXAMPLE".into(),
150− secret_access_key: "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY".into(),
151− region: "us-east-1".into(),
152− }
153− }
154−
155− /// AWS's own example of a presigned URL (Authenticating Requests:
156− /// Using Query Parameters).
157− #[test]
158− fn a_presigned_url_matches_the_aws_example() {
159− let url = example().presign_get(
160− "https://examplebucket.s3.amazonaws.com",
161− "examplebucket.s3.amazonaws.com",
162− "/test.txt",
163− "20130524T000000Z",
164− 86400,
165− );
166− assert!(url.starts_with("https://examplebucket.s3.amazonaws.com/test.txt?X-Amz-Algorithm=AWS4-HMAC-SHA256"));
167− assert!(url.contains("X-Amz-Credential=AKIAIOSFODNN7EXAMPLE%2F20130524%2Fus-east-1%2Fs3%2Faws4_request"));
168− assert!(url.contains("&X-Amz-Signature=aeeed9bbccd4d02ee5c0109b86d86835f995330da4c265957d157751f604d404&"), "{url}");
169− }
170−
171− /// AWS's own example of a signed GET with a range (Authenticating
172− /// Requests: Using the Authorization Header).
173− #[test]
174− fn a_signed_request_matches_the_aws_example() {
175− let empty = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855";
176− let headers = [
177− ("Host", "examplebucket.s3.amazonaws.com"),
178− ("Range", "bytes=0-9"),
179− ("x-amz-content-sha256", empty),
180− ("x-amz-date", "20130524T000000Z"),
181− ]
182− .map(|(name, value)| (name.to_owned(), value.to_owned()));
183− let authorization = example().authorization("GET", "/test.txt", &[], &headers, empty);
184− assert_eq!(
185− authorization,
186− "AWS4-HMAC-SHA256 Credential=AKIAIOSFODNN7EXAMPLE/20130524/us-east-1/s3/aws4_request, \
187− SignedHeaders=host;range;x-amz-content-sha256;x-amz-date, \
188− Signature=f0e8bdb87c964420e857bd35b5d6ed310bd44f0170aba48dd91039c6036bdb41"
189− );
190− }
191−
192− #[test]
193− fn dates_and_encoding() {
194− assert_eq!(amz_date(1_369_353_600_000), "20130524T000000Z");
195− assert_eq!(uri_encode("a b/c+d~", true), "a%20b/c%2Bd~");
196− assert_eq!(uri_encode("a/b", false), "a%2Fb");
197− let query = [("uploadId".to_owned(), "x y".to_owned()), ("partNumber".to_owned(), "2".to_owned())];
198− assert_eq!(canonical_query(&query), "partNumber=2&uploadId=x%20y");
199− assert_eq!(canonical_query(&[("uploads".to_owned(), String::new())]), "uploads=");
200− }
201−}
+32−0
1+//! Where packages' files are kept: the `BlobStore` port (crates/blobstore),
2+//! with R2 behind it on Cloudflare and any S3-compatible storage (MinIO in
3+//! the compose file) when self-hosted. BLOB_STORE chooses: `r2` (the
4+//! default) or `s3`.
5+//!
6+//! Files are content-addressed: a blob stored whole is at
7+//! `blobs/sha256/<hex>`, and one that came in parts at the key its upload
8+//! started with, which the `blobs` table records. Large uploads go up as
9+//! multipart parts of one size, as R2 requires (every part but the last
10+//! the same size), however the client cut its chunks.
11+
12+use worker::{Env, Result};
13+
14+#[cfg(test)]
15+pub use g1t_blobstore::Got;
16+pub use g1t_blobstore::{BlobStore, Part, Store, var};
17+
18+/// The `BLOBS` bucket; R2's S3 endpoint signs downloads when
19+/// R2_ACCESS_KEY_ID, R2_SECRET_ACCESS_KEY, R2_ACCOUNT_ID and R2_BUCKET are
20+/// set. Self-hosted: S3_BUCKET, and S3_PUBLIC_ENDPOINT for signed downloads.
21+const CONFIG: g1t_blobstore::Config = g1t_blobstore::Config {
22+ kind: "BLOB_STORE",
23+ binding: "BLOBS",
24+ r2_signer: Some(["R2_ACCESS_KEY_ID", "R2_SECRET_ACCESS_KEY", "R2_ACCOUNT_ID", "R2_BUCKET"]),
25+ s3_bucket: "S3_BUCKET",
26+ s3_public_endpoint: Some("S3_PUBLIC_ENDPOINT"),
27+};
28+
29+/// The store this installation keeps packages' files in.
30+pub fn from_env(env: &Env) -> Result<Store> {
31+ Store::from_env(env, &CONFIG)
32+}
+0−159
1−//! Where packages' files are kept: the `BlobStore` port, with R2 behind it
2−//! on Cloudflare and any S3-compatible storage (MinIO in the compose file)
3−//! when self-hosted. BLOB_STORE chooses: `r2` (the default) or `s3`.
4−//!
5−//! Files are content-addressed: a blob stored whole is at
6−//! `blobs/sha256/<hex>`, and one that came in parts at the key its upload
7−//! started with, which the `blobs` table records. Large uploads go up as
8−//! multipart parts of one size, as R2 requires (every part but the last
9−//! the same size), however the client cut its chunks.
10−
11−mod r2;
12−mod s3;
13−
14−use serde::{Deserialize, Serialize};
15−use worker::{Env, Response, ResponseBody, Result};
16−
17−use crate::range::Wanted;
18−
19−pub use r2::R2Store;
20−pub use s3::S3Store;
21−
22−/// One part of a multipart upload, as completing it needs.
23−#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
24−pub struct Part {
25− pub number: u16,
26− pub etag: String,
27−}
28−
29−/// An object read back.
30−pub struct Got {
31− /// The whole object's size, whatever range was read.
32− pub size: u64,
33− pub body: ResponseBody,
34−}
35−
36−impl Got {
37− pub async fn bytes(self) -> Result<Vec<u8>> {
38− match self.body {
39− ResponseBody::Empty => Ok(Vec::new()),
40− ResponseBody::Body(bytes) => Ok(bytes),
41− stream => Response::from_body(stream)?.bytes().await,
42− }
43− }
44−}
45−
46−/// What the registry needs of storage.
47−pub trait BlobStore {
48− async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()>;
49− /// The object, or the part of it `range` asks for.
50− async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>>;
51− /// The object's size, if it is there.
52− async fn head(&self, key: &str) -> Result<Option<u64>>;
53− async fn delete(&self, key: &str) -> Result<()>;
54− /// Starts a multipart upload to `key`, and says its id.
55− async fn create_multipart(&self, key: &str) -> Result<String>;
56− async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part>;
57− async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()>;
58− async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()>;
59− /// A URL that downloads the object for `expires` seconds without
60− /// passing through this Worker, when the store can sign one.
61− fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String>;
62−
63− /// The whole object, read into memory: for small ones only.
64− async fn read(&self, key: &str) -> Result<Option<Vec<u8>>> {
65− match self.get(key, None).await? {
66− Some(got) => Ok(Some(got.bytes().await?)),
67− None => Ok(None),
68− }
69− }
70−}
71−
72−/// The store this installation is configured with.
73−pub enum Store {
74− R2(R2Store),
75− S3(S3Store),
76−}
77−
78−impl Store {
79− pub fn from_env(env: &Env) -> Result<Store> {
80− let kind = env.var("BLOB_STORE").map(|v| v.to_string()).unwrap_or_default();
81− if kind == "s3" {
82− return Ok(Store::S3(S3Store::from_env(env)?));
83− }
84− Ok(Store::R2(R2Store::from_env(env)?))
85− }
86−}
87−
88−impl BlobStore for Store {
89− async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
90− match self {
91− Store::R2(s) => s.put(key, bytes).await,
92− Store::S3(s) => s.put(key, bytes).await,
93− }
94− }
95−
96− async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>> {
97− match self {
98− Store::R2(s) => s.get(key, range).await,
99− Store::S3(s) => s.get(key, range).await,
100− }
101− }
102−
103− async fn head(&self, key: &str) -> Result<Option<u64>> {
104− match self {
105− Store::R2(s) => s.head(key).await,
106− Store::S3(s) => s.head(key).await,
107− }
108− }
109−
110− async fn delete(&self, key: &str) -> Result<()> {
111− match self {
112− Store::R2(s) => s.delete(key).await,
113− Store::S3(s) => s.delete(key).await,
114− }
115− }
116−
117− async fn create_multipart(&self, key: &str) -> Result<String> {
118− match self {
119− Store::R2(s) => s.create_multipart(key).await,
120− Store::S3(s) => s.create_multipart(key).await,
121− }
122− }
123−
124− async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part> {
125− match self {
126− Store::R2(s) => s.upload_part(key, upload_id, number, bytes).await,
127− Store::S3(s) => s.upload_part(key, upload_id, number, bytes).await,
128− }
129− }
130−
131− async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()> {
132− match self {
133− Store::R2(s) => s.complete_multipart(key, upload_id, parts).await,
134− Store::S3(s) => s.complete_multipart(key, upload_id, parts).await,
135− }
136− }
137−
138− async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()> {
139− match self {
140− Store::R2(s) => s.abort_multipart(key, upload_id).await,
141− Store::S3(s) => s.abort_multipart(key, upload_id).await,
142− }
143− }
144−
145− fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String> {
146− match self {
147− Store::R2(s) => s.presign_get(key, expires, now_ms),
148− Store::S3(s) => s.presign_get(key, expires, now_ms),
149− }
150− }
151−}
152−
153−/// A configuration variable, or the empty string.
154−pub(crate) fn var(env: &Env, name: &str) -> String {
155− env.var(name)
156− .map(|v| v.to_string())
157− .or_else(|_| env.secret(name).map(|v| v.to_string()))
158− .unwrap_or_default()
159−}
+0−93
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−}
+0−249
1−//! The S3 adapter, for self-hosted installations: any S3-compatible store
2−//! (MinIO, Ceph, Garage, AWS) over fetch, signed with SigV4, path-style.
3−//! S3_ENDPOINT, S3_BUCKET, S3_ACCESS_KEY_ID, S3_SECRET_ACCESS_KEY and
4−//! S3_REGION say where; S3_PUBLIC_ENDPOINT, when set, is the address
5−//! clients reach the store at, and large downloads are then sent there
6−//! with a signed URL instead of through the Worker.
7−
8−use worker::wasm_bindgen::JsValue;
9−use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Response, Result, Url};
10−
11−use super::{BlobStore, Got, Part, var};
12−use crate::range::Wanted;
13−use crate::sigv4::{Credentials, UNSIGNED, amz_date};
14−
15−pub struct S3Store {
16− /// `http://minio:9000`, without a trailing slash.
17− endpoint: String,
18− /// Where clients reach the same store, for signed URLs.
19− public_endpoint: Option<String>,
20− bucket: String,
21− credentials: Credentials,
22−}
23−
24−fn failed(what: &str, status: u16, body: &str) -> worker::Error {
25− let said: String = body.chars().take(300).collect();
26− worker::Error::RustError(format!("storage {what} failed with status {status}: {said}"))
27−}
28−
29−/// The text of the first `<tag>` in an XML answer.
30−fn xml_value<'a>(xml: &'a str, tag: &str) -> Option<&'a str> {
31− let open = format!("<{tag}>");
32− let start = xml.find(&open)? + open.len();
33− let end = xml[start..].find(&format!("</{tag}>"))? + start;
34− Some(&xml[start..end])
35−}
36−
37−fn host_of(endpoint: &str) -> String {
38− endpoint
39− .split_once("://")
40− .map_or(endpoint, |(_, rest)| rest)
41− .split('/')
42− .next()
43− .unwrap_or_default()
44− .to_owned()
45−}
46−
47−impl S3Store {
48− pub fn from_env(env: &Env) -> Result<S3Store> {
49− let endpoint = var(env, "S3_ENDPOINT").trim_end_matches('/').to_owned();
50− let bucket = var(env, "S3_BUCKET");
51− if endpoint.is_empty() || bucket.is_empty() {
52− return Err(worker::Error::RustError("BLOB_STORE is s3, but S3_ENDPOINT or S3_BUCKET is not set".into()));
53− }
54− let region = var(env, "S3_REGION");
55− let public = var(env, "S3_PUBLIC_ENDPOINT").trim_end_matches('/').to_owned();
56− Ok(S3Store {
57− endpoint,
58− public_endpoint: (!public.is_empty()).then_some(public),
59− bucket,
60− credentials: Credentials {
61− access_key_id: var(env, "S3_ACCESS_KEY_ID"),
62− secret_access_key: var(env, "S3_SECRET_ACCESS_KEY"),
63− region: if region.is_empty() { "us-east-1".to_owned() } else { region },
64− },
65− })
66− }
67−
68− fn path(&self, key: &str) -> String {
69− format!("/{}/{key}", self.bucket)
70− }
71−
72− /// Sends one signed request, and answers with the response whatever
73− /// its status.
74− async fn send(
75− &self,
76− method: Method,
77− key: &str,
78− query: &[(String, String)],
79− extra: &[(&str, String)],
80− body: Option<Vec<u8>>,
81− ) -> Result<Response> {
82− let path = self.path(key);
83− let date = amz_date(g1t_kit::now_ms());
84− let mut signed = vec![
85− ("host".to_owned(), host_of(&self.endpoint)),
86− ("x-amz-content-sha256".to_owned(), UNSIGNED.to_owned()),
87− ("x-amz-date".to_owned(), date),
88− ];
89− for (name, value) in extra {
90− signed.push(((*name).to_owned(), value.clone()));
91− }
92− let authorization = self
93− .credentials
94− .authorization(method.as_ref(), &path, query, &signed, UNSIGNED);
95− let headers = Headers::new();
96− for (name, value) in &signed {
97− if name != "host" {
98− headers.set(name, value)?;
99− }
100− }
101− headers.set("authorization", &authorization)?;
102− let mut url = Url::parse(&format!("{}{}", self.endpoint, crate::sigv4::uri_encode(&path, true)))?;
103− if !query.is_empty() {
104− let text: Vec<String> = query
105− .iter()
106− .map(|(k, v)| {
107− let (k, v) = (crate::sigv4::uri_encode(k, false), crate::sigv4::uri_encode(v, false));
108− if v.is_empty() { format!("{k}=") } else { format!("{k}={v}") }
109− })
110− .collect();
111− url.set_query(Some(&text.join("&")));
112− }
113− let mut init = RequestInit::new();
114− init.with_method(method).with_headers(headers);
115− if let Some(body) = body {
116− init.with_body(Some(JsValue::from(worker::js_sys::Uint8Array::from(body.as_slice()))));
117− }
118− Fetch::Request(Request::new_with_init(url.as_str(), &init)?).send().await
119− }
120−
121− async fn ok(&self, what: &str, mut response: Response) -> Result<Response> {
122− let status = response.status_code();
123− if (200..300).contains(&status) {
124− return Ok(response);
125− }
126− let body = response.text().await.unwrap_or_default();
127− Err(failed(what, status, &body))
128− }
129−}
130−
131−fn query(pairs: &[(&str, &str)]) -> Vec<(String, String)> {
132− pairs.iter().map(|(k, v)| ((*k).to_owned(), (*v).to_owned())).collect()
133−}
134−
135−impl BlobStore for S3Store {
136− async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
137− let length = bytes.len().to_string();
138− let response = self
139− .send(Method::Put, key, &[], &[("content-length", length)], Some(bytes))
140− .await?;
141− self.ok("put", response).await.map(|_| ())
142− }
143−
144− async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>> {
145− let extra: Vec<(&str, String)> = range
146− .map(|r| ("range", format!("bytes={}-{}", r.offset, r.offset + r.length - 1)))
147− .into_iter()
148− .collect();
149− let response = self.send(Method::Get, key, &[], &extra, None).await?;
150− if response.status_code() == 404 {
151− return Ok(None);
152− }
153− let response = self.ok("get", response).await?;
154− let size = match response.headers().get("content-range")? {
155− // `bytes 0-9/100`: the whole object's size is after the slash.
156− Some(range) => range.rsplit('/').next().and_then(|n| n.parse().ok()).unwrap_or(0),
157− None => response.headers().get("content-length")?.and_then(|n| n.parse().ok()).unwrap_or(0),
158− };
159− let (_, body) = response.into_parts();
160− Ok(Some(Got { size, body }))
161− }
162−
163− async fn head(&self, key: &str) -> Result<Option<u64>> {
164− let response = self.send(Method::Head, key, &[], &[], None).await?;
165− if response.status_code() == 404 {
166− return Ok(None);
167− }
168− let response = self.ok("head", response).await?;
169− Ok(response.headers().get("content-length")?.and_then(|n| n.parse().ok()))
170− }
171−
172− async fn delete(&self, key: &str) -> Result<()> {
173− let response = self.send(Method::Delete, key, &[], &[], None).await?;
174− if response.status_code() == 404 {
175− return Ok(());
176− }
177− self.ok("delete", response).await.map(|_| ())
178− }
179−
180− async fn create_multipart(&self, key: &str) -> Result<String> {
181− let response = self.send(Method::Post, key, &query(&[("uploads", "")]), &[], None).await?;
182− let text = self.ok("create multipart", response).await?.text().await?;
183− xml_value(&text, "UploadId")
184− .map(str::to_owned)
185− .ok_or_else(|| failed("create multipart", 200, &text))
186− }
187−
188− async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part> {
189− let number_text = number.to_string();
190− let length = bytes.len().to_string();
191− let response = self
192− .send(
193− Method::Put,
194− key,
195− &query(&[("partNumber", &number_text), ("uploadId", upload_id)]),
196− &[("content-length", length)],
197− Some(bytes),
198− )
199− .await?;
200− let response = self.ok("upload part", response).await?;
201− let etag = response.headers().get("etag")?.unwrap_or_default();
202− Ok(Part { number, etag })
203− }
204−
205− async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()> {
206− let mut xml = String::from("<CompleteMultipartUpload>");
207− for part in parts {
208− xml.push_str(&format!("<Part><PartNumber>{}</PartNumber><ETag>{}</ETag></Part>", part.number, part.etag));
209− }
210− xml.push_str("</CompleteMultipartUpload>");
211− let length = xml.len().to_string();
212− let response = self
213− .send(Method::Post, key, &query(&[("uploadId", upload_id)]), &[("content-length", length)], Some(xml.into_bytes()))
214− .await?;
215− // S3 may answer 200 and still have failed, saying so in the body.
216− let text = self.ok("complete multipart", response).await?.text().await?;
217− if text.contains("<Error>") {
218− return Err(failed("complete multipart", 200, &text));
219− }
220− Ok(())
221− }
222−
223− async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()> {
224− let response = self.send(Method::Delete, key, &query(&[("uploadId", upload_id)]), &[], None).await?;
225− if response.status_code() == 404 {
226− return Ok(());
227− }
228− self.ok("abort multipart", response).await.map(|_| ())
229− }
230−
231− fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String> {
232− let base = self.public_endpoint.as_ref()?;
233− Some(self.credentials.presign_get(base, &host_of(base), &self.path(key), &amz_date(now_ms), expires))
234− }
235−}
236−
237−#[cfg(test)]
238−mod tests {
239− use super::*;
240−
241− #[test]
242− fn answers_are_read_from_their_xml() {
243− let xml = "<InitiateMultipartUploadResult><Bucket>b</Bucket><UploadId>abc-123</UploadId></InitiateMultipartUploadResult>";
244− assert_eq!(xml_value(xml, "UploadId"), Some("abc-123"));
245− assert_eq!(xml_value(xml, "Key"), None);
246− assert_eq!(host_of("http://minio:9000"), "minio:9000");
247− assert_eq!(host_of("https://s3.example.com/base"), "s3.example.com");
248− }
249−}