g1t/services/packages/src/store/mod.rs

159 lines5,515 bytesCodeBlame
1//! Where packages' files are kept: the `BlobStore` port, with R2 behind it
2//! on Cloudflare and any S3-compatible storage (MinIO in the compose file)
3//! when self-hosted. BLOB_STORE chooses: `r2` (the default) or `s3`.
4//!
5//! Files are content-addressed: a blob stored whole is at
6//! `blobs/sha256/<hex>`, and one that came in parts at the key its upload
7//! started with, which the `blobs` table records. Large uploads go up as
8//! multipart parts of one size, as R2 requires (every part but the last
9//! the same size), however the client cut its chunks.
10
11mod r2;
12mod s3;
13
14use serde::{Deserialize, Serialize};
15use worker::{Env, Response, ResponseBody, Result};
16
17use crate::range::Wanted;
18
19pub use r2::R2Store;
20pub use s3::S3Store;
21
22/// One part of a multipart upload, as completing it needs.
23#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
24pub struct Part {
25 pub number: u16,
26 pub etag: String,
27}
28
29/// An object read back.
30pub struct Got {
31 /// The whole object's size, whatever range was read.
32 pub size: u64,
33 pub body: ResponseBody,
34}
35
36impl Got {
37 pub async fn bytes(self) -> Result<Vec<u8>> {
38 match self.body {
39 ResponseBody::Empty => Ok(Vec::new()),
40 ResponseBody::Body(bytes) => Ok(bytes),
41 stream => Response::from_body(stream)?.bytes().await,
42 }
43 }
44}
45
46/// What the registry needs of storage.
47pub trait BlobStore {
48 async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()>;
49 /// The object, or the part of it `range` asks for.
50 async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>>;
51 /// The object's size, if it is there.
52 async fn head(&self, key: &str) -> Result<Option<u64>>;
53 async fn delete(&self, key: &str) -> Result<()>;
54 /// Starts a multipart upload to `key`, and says its id.
55 async fn create_multipart(&self, key: &str) -> Result<String>;
56 async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part>;
57 async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()>;
58 async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()>;
59 /// A URL that downloads the object for `expires` seconds without
60 /// passing through this Worker, when the store can sign one.
61 fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String>;
62
63 /// The whole object, read into memory: for small ones only.
64 async fn read(&self, key: &str) -> Result<Option<Vec<u8>>> {
65 match self.get(key, None).await? {
66 Some(got) => Ok(Some(got.bytes().await?)),
67 None => Ok(None),
68 }
69 }
70}
71
72/// The store this installation is configured with.
73pub enum Store {
74 R2(R2Store),
75 S3(S3Store),
76}
77
78impl Store {
79 pub fn from_env(env: &Env) -> Result<Store> {
80 let kind = env.var("BLOB_STORE").map(|v| v.to_string()).unwrap_or_default();
81 if kind == "s3" {
82 return Ok(Store::S3(S3Store::from_env(env)?));
83 }
84 Ok(Store::R2(R2Store::from_env(env)?))
85 }
86}
87
88impl BlobStore for Store {
89 async fn put(&self, key: &str, bytes: Vec<u8>) -> Result<()> {
90 match self {
91 Store::R2(s) => s.put(key, bytes).await,
92 Store::S3(s) => s.put(key, bytes).await,
93 }
94 }
95
96 async fn get(&self, key: &str, range: Option<Wanted>) -> Result<Option<Got>> {
97 match self {
98 Store::R2(s) => s.get(key, range).await,
99 Store::S3(s) => s.get(key, range).await,
100 }
101 }
102
103 async fn head(&self, key: &str) -> Result<Option<u64>> {
104 match self {
105 Store::R2(s) => s.head(key).await,
106 Store::S3(s) => s.head(key).await,
107 }
108 }
109
110 async fn delete(&self, key: &str) -> Result<()> {
111 match self {
112 Store::R2(s) => s.delete(key).await,
113 Store::S3(s) => s.delete(key).await,
114 }
115 }
116
117 async fn create_multipart(&self, key: &str) -> Result<String> {
118 match self {
119 Store::R2(s) => s.create_multipart(key).await,
120 Store::S3(s) => s.create_multipart(key).await,
121 }
122 }
123
124 async fn upload_part(&self, key: &str, upload_id: &str, number: u16, bytes: Vec<u8>) -> Result<Part> {
125 match self {
126 Store::R2(s) => s.upload_part(key, upload_id, number, bytes).await,
127 Store::S3(s) => s.upload_part(key, upload_id, number, bytes).await,
128 }
129 }
130
131 async fn complete_multipart(&self, key: &str, upload_id: &str, parts: &[Part]) -> Result<()> {
132 match self {
133 Store::R2(s) => s.complete_multipart(key, upload_id, parts).await,
134 Store::S3(s) => s.complete_multipart(key, upload_id, parts).await,
135 }
136 }
137
138 async fn abort_multipart(&self, key: &str, upload_id: &str) -> Result<()> {
139 match self {
140 Store::R2(s) => s.abort_multipart(key, upload_id).await,
141 Store::S3(s) => s.abort_multipart(key, upload_id).await,
142 }
143 }
144
145 fn presign_get(&self, key: &str, expires: u32, now_ms: u64) -> Option<String> {
146 match self {
147 Store::R2(s) => s.presign_get(key, expires, now_ms),
148 Store::S3(s) => s.presign_get(key, expires, now_ms),
149 }
150 }
151}
152
153/// A configuration variable, or the empty string.
154pub(crate) fn var(env: &Env, name: &str) -> String {
155 env.var(name)
156 .map(|v| v.to_string())
157 .or_else(|_| env.secret(name).map(|v| v.to_string()))
158 .unwrap_or_default()
159}