g1t/crates/blobstore/src/lib.rs

192 lines6,664 bytesCodeBlame
1//! Object storage behind one port, `BlobStore`: R2 on Cloudflare, and any
2//! S3-compatible store (MinIO in the self-host compose file) elsewhere.
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
13mod r2;
14mod s3;
15pub mod sigv4;
16
17use serde::{Deserialize, Serialize};
18use worker::{Env, Response, ResponseBody, Result};
19
20pub use r2::R2Store;
21pub 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)]
25pub struct Wanted {
26 pub offset: u64,
27 pub length: u64,
28}
29
30impl 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)]
39pub struct Part {
40 pub number: u16,
41 pub etag: String,
42}
43
44/// An object read back.
45pub struct Got {
46 /// The whole object's size, whatever range was read.
47 pub size: u64,
48 pub body: ResponseBody,
49}
50
51impl 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)]
63pub 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)]
91pub 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.
107pub enum Store {
108 R2(R2Store),
109 S3(S3Store),
110}
111
112impl 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
121impl 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.
187pub 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}