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

249 lines9,982 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 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
8use worker::wasm_bindgen::JsValue;
9use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Response, Result, Url};
10
11use super::{BlobStore, Got, Part, var};
12use crate::range::Wanted;
13use crate::sigv4::{Credentials, UNSIGNED, amz_date};
14
15pub 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
24fn 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.
30fn 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
37fn 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
47impl 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
131fn query(pairs: &[(&str, &str)]) -> Vec<(String, String)> {
132 pairs.iter().map(|(k, v)| ((*k).to_owned(), (*v).to_owned())).collect()
133}
134
135impl 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)]
238mod 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}