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 | ||
| 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.