g1t/crates/blobstore/src/s3.rs

255 lines10,234 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.

Merge branch 'worktree-agent-ac5b181a013e54348'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
9use worker::wasm_bindgen::JsValue;
10use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Response, Result, Url};
11
12use crate::sigv4::{Credentials, UNSIGNED, amz_date};
13use crate::{BlobStore, Config, Got, Part, Wanted, var};
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, 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
137fn query(pairs: &[(&str, &str)]) -> Vec<(String, String)> {
138 pairs.iter().map(|(k, v)| ((*k).to_owned(), (*v).to_owned())).collect()
139}
140
141impl 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)]
244mod 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}

This file's history is long; its oldest lines are credited to the oldest commit read.