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 | //! Object storage behind one port, `BlobStore`: R2 on Cloudflare, and any |
| Merge branch 'worktree-agent-af58ac8933b0dd125' | 2 | //! S3-compatible store (RustFS in the self-host compose file) elsewhere. |
| Merge branch 'worktree-agent-ac5b181a013e54348' | 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.