| 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 | } |