pr_01m47d24b0e6n91zwymwxg0vpx/services/repos/src/store.rs

295 lines10,088 bytesCodeBlame
1//! The storage that actually holds git repositories.
2//!
3//! The service depends on the [`GitStore`] and [`GitRepo`] ports;
4//! [`ArtifactsStore`] is the adapter for Cloudflare Artifacts.
5
6use g1t_contracts::repos::{Branch, Commit, EntryKind, GitAccess, Signature, TreeEntry};
7use g1t_contracts::time::rfc3339;
8use g1t_kit::js;
9use serde::Deserialize;
10use worker::js_sys::{Reflect, Uint8Array};
11use worker::wasm_bindgen::{JsCast, JsValue};
12use worker::{Env, Result};
13
14/// How long a credential handed to git stays valid.
15const TOKEN_TTL_SECONDS: u32 = 300;
16
17#[derive(Clone, Copy)]
18pub enum Scope {
19 Read,
20 Write,
21}
22
23/// A place repositories live. `key` is the store's own name for a repo.
24#[allow(async_fn_in_trait)]
25pub trait GitStore {
26 type Repo: GitRepo;
27
28 /// Creates an empty repository. Succeeds if it already exists.
29 async fn create(
30 &self,
31 key: &str,
32 description: Option<&str>,
33 default_branch: &str,
34 ) -> Result<()>;
35 async fn open(&self, key: &str) -> Result<Self::Repo>;
36 /// Removes a repository and everything in it, for good. Succeeds if it
37 /// is already gone.
38 async fn delete(&self, key: &str) -> Result<()>;
39}
40
41/// One open repository.
42#[allow(async_fn_in_trait)]
43pub trait GitRepo {
44 /// A remote URL and short-lived credential for git itself.
45 async fn access(&self, scope: Scope) -> Result<GitAccess>;
46 /// Every branch and the commit it points to.
47 async fn branches(&self) -> Result<Vec<Branch>>;
48 /// Newest first along the first-parent chain; empty for an unknown ref.
49 async fn log(&self, git_ref: &str, limit: u32) -> Result<Vec<Commit>>;
50 /// The parents of a commit, or `None` if the commit does not exist.
51 async fn parents(&self, commit_hash: &str) -> Result<Option<Vec<String>>>;
52 async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>>;
53 async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>>;
54 /// `None` when the ref or path does not resolve to a file.
55 async fn read_file(&self, git_ref: &str, path: &str) -> Result<Option<Vec<u8>>>;
56 /// Makes a copy-on-write copy of this repository under `target_key`.
57 async fn fork(&self, target_key: &str) -> Result<()>;
58}
59
60pub struct ArtifactsStore {
61 binding: JsValue,
62}
63
64impl ArtifactsStore {
65 pub fn new(env: &Env) -> Result<Self> {
66 Ok(Self {
67 binding: js::binding(env, "ARTIFACTS")?,
68 })
69 }
70}
71
72impl GitStore for ArtifactsStore {
73 type Repo = ArtifactsRepo;
74
75 async fn create(
76 &self,
77 key: &str,
78 description: Option<&str>,
79 default_branch: &str,
80 ) -> Result<()> {
81 let options = js::to_js(&serde_json::json!({
82 "description": description,
83 "setDefaultBranch": default_branch,
84 }))?;
85 match js::call(&self.binding, "create", &[key.into(), options]).await {
86 // Left behind by an earlier failed attempt; adopt it.
87 Err(thrown) if !thrown.is("ALREADY_EXISTS") => Err(thrown.into()),
88 _ => Ok(()),
89 }
90 }
91
92 async fn delete(&self, key: &str) -> Result<()> {
93 match js::call(&self.binding, "delete", &[key.into()]).await {
94 // Gone already: an earlier purge got this far.
95 Err(thrown) if !thrown.is("NOT_FOUND") => Err(thrown.into()),
96 _ => Ok(()),
97 }
98 }
99
100 async fn open(&self, key: &str) -> Result<ArtifactsRepo> {
101 Ok(ArtifactsRepo {
102 handle: js::call(&self.binding, "get", &[key.into()]).await?,
103 key: key.to_owned(),
104 })
105 }
106}
107
108/// A handle to one Artifacts repository. It is an RPC stub, so it is
109/// released when dropped.
110pub struct ArtifactsRepo {
111 handle: JsValue,
112 /// The repository's store key, which scopes its cached objects.
113 key: String,
114}
115
116/// Where cached git objects live. Trees and blobs are named by their
117/// content, so a cached one is never stale; each is kept under its own
118/// repository's key, so a repository only ever finds its own objects.
119const OBJECT_CACHE: &str = "https://objects.g1t.internal/";
120/// Blobs larger than this are not cached.
121const MAX_CACHED_BLOB: usize = 1024 * 1024;
122const OBJECT_MAX_AGE: &str = "public, max-age=31536000, immutable";
123
124impl ArtifactsRepo {
125 fn cache_url(&self, kind: &str, hash: &str) -> String {
126 format!("{OBJECT_CACHE}{}/{kind}/{hash}", self.key)
127 }
128
129 async fn cached(&self, kind: &str, hash: &str) -> Option<Vec<u8>> {
130 let mut response = worker::Cache::default()
131 .get(self.cache_url(kind, hash), false)
132 .await
133 .ok()??;
134 response.bytes().await.ok()
135 }
136
137 /// Keeps an object for next time. A failure only costs a later read.
138 async fn keep(&self, kind: &str, hash: &str, bytes: Vec<u8>) {
139 let Ok(mut response) = worker::Response::from_bytes(bytes) else {
140 return;
141 };
142 let _ = response.headers_mut().set("cache-control", OBJECT_MAX_AGE);
143 let _ = worker::Cache::default()
144 .put(self.cache_url(kind, hash), response)
145 .await;
146 }
147}
148
149impl Drop for ArtifactsRepo {
150 fn drop(&mut self) {
151 let symbol = js::get(&worker::js_sys::global(), "Symbol");
152 let dispose = js::get(&symbol, "dispose");
153 if let Ok(function) = Reflect::get(&self.handle, &dispose)
154 .and_then(|value| value.dyn_into::<worker::js_sys::Function>())
155 {
156 let _ = function.call0(&self.handle);
157 }
158 }
159}
160
161#[derive(Deserialize)]
162#[serde(rename_all = "camelCase")]
163struct RawCommit {
164 hash: String,
165 tree_hash: String,
166 message: String,
167 author: Signature,
168 parents: Vec<String>,
169 /// Seconds since the epoch.
170 authored_at: u64,
171}
172
173#[derive(Deserialize)]
174struct RawEntry {
175 name: String,
176 hash: String,
177 #[serde(rename = "type")]
178 kind: EntryKind,
179}
180
181#[derive(Deserialize)]
182struct RawInfo {
183 remote: String,
184}
185
186#[derive(Deserialize)]
187struct RawToken {
188 plaintext: String,
189}
190
191/// The bytes of a `Blob`, or `None` for null.
192async fn blob_bytes(blob: JsValue) -> Result<Option<Vec<u8>>> {
193 if blob.is_null() || blob.is_undefined() {
194 return Ok(None);
195 }
196 let buffer = js::call(&blob, "arrayBuffer", &[]).await?;
197 Ok(Some(Uint8Array::new(&buffer).to_vec()))
198}
199
200impl GitRepo for ArtifactsRepo {
201 async fn access(&self, scope: Scope) -> Result<GitAccess> {
202 let scope = match scope {
203 Scope::Read => "read",
204 Scope::Write => "write",
205 };
206 let info: RawInfo = js::from_js(&js::call(&self.handle, "info", &[]).await?)?;
207 let token: RawToken = js::from_js(
208 &js::call(
209 &self.handle,
210 "createToken",
211 &[scope.into(), TOKEN_TTL_SECONDS.into()],
212 )
213 .await?,
214 )?;
215 Ok(GitAccess {
216 remote: info.remote,
217 token: token.plaintext,
218 })
219 }
220
221 async fn branches(&self) -> Result<Vec<Branch>> {
222 crate::refs::branches(&self.access(Scope::Read).await?).await
223 }
224
225 async fn log(&self, git_ref: &str, limit: u32) -> Result<Vec<Commit>> {
226 let options = js::to_js(&serde_json::json!({ "ref": git_ref, "limit": limit }))?;
227 let commits: Vec<RawCommit> =
228 js::from_js(&js::call(&self.handle, "log", &[options]).await?)?;
229 Ok(commits
230 .into_iter()
231 .map(|commit| Commit {
232 hash: commit.hash,
233 tree_hash: commit.tree_hash,
234 message: commit.message,
235 author: commit.author,
236 parents: commit.parents,
237 authored_at: rfc3339(commit.authored_at * 1000),
238 })
239 .collect())
240 }
241
242 async fn parents(&self, commit_hash: &str) -> Result<Option<Vec<String>>> {
243 let commit: Option<RawCommit> =
244 js::from_js(&js::call(&self.handle, "readCommit", &[commit_hash.into()]).await?)?;
245 Ok(commit.map(|commit| commit.parents))
246 }
247
248 async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>> {
249 if let Some(bytes) = self.cached("tree", tree_hash).await
250 && let Ok(entries) = serde_json::from_slice::<Vec<TreeEntry>>(&bytes) {
251 return Ok(Some(entries));
252 }
253 let entries: Option<Vec<RawEntry>> =
254 js::from_js(&js::call(&self.handle, "readTree", &[tree_hash.into()]).await?)?;
255 let entries: Option<Vec<TreeEntry>> = entries.map(|entries| {
256 entries
257 .into_iter()
258 .map(|entry| TreeEntry {
259 name: entry.name,
260 hash: entry.hash,
261 kind: entry.kind,
262 })
263 .collect()
264 });
265 if let Some(entries) = &entries
266 && let Ok(bytes) = serde_json::to_vec(entries) {
267 self.keep("tree", tree_hash, bytes).await;
268 }
269 Ok(entries)
270 }
271
272 async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>> {
273 if let Some(bytes) = self.cached("blob", blob_hash).await {
274 return Ok(Some(bytes));
275 }
276 let bytes = blob_bytes(js::call(&self.handle, "readBlob", &[blob_hash.into()]).await?).await?;
277 if let Some(bytes) = bytes.as_ref().filter(|bytes| bytes.len() <= MAX_CACHED_BLOB) {
278 self.keep("blob", blob_hash, bytes.clone()).await;
279 }
280 Ok(bytes)
281 }
282
283 async fn read_file(&self, git_ref: &str, path: &str) -> Result<Option<Vec<u8>>> {
284 let args = js::to_js(&serde_json::json!({ "ref": git_ref, "path": path }))?;
285 blob_bytes(js::call(&self.handle, "readFile", &[args]).await?).await
286 }
287
288 async fn fork(&self, target_key: &str) -> Result<()> {
289 let options = js::to_js(&serde_json::json!({ "defaultBranchOnly": true }))?;
290 match js::call(&self.handle, "fork", &[target_key.into(), options]).await {
291 Err(thrown) if !thrown.is("ALREADY_EXISTS") => Err(thrown.into()),
292 _ => Ok(()),
293 }
294 }
295}