pr_01m47d15m3e54sn21z27rpy5n9/apps/api/src/blobs.rs

257 lines10,308 bytesCodeBlame
1//! Artifacts and the cache of GitHub Actions jobs, kept in Workers KV in
2//! chunks, with KV's own expiry: artifacts for 14 days with their run,
3//! cache entries for 7 days with their repository.
4//!
5//! A sandbox reaches these with its job's token, at
6//! `/actions/jobs/{job}/artifacts[/{name}]` and `/actions/jobs/{job}/cache`.
7//! People download an artifact at
8//! `/repos/{owner}/{repo}/actions/runs/{run}/artifacts/{name}`.
9
10use serde::{Deserialize, Serialize};
11use serde_json::{Value, json};
12use worker::kv::KvStore;
13use worker::{Env, Request, Response, Result};
14
15use g1t_contracts::{FailureCode, Outcome};
16
17use crate::operations::Services;
18
19/// KV's largest value is 25 MiB; chunks stay under it.
20const CHUNK: usize = 20 * 1024 * 1024;
21/// The largest artifact or cache entry, kept within a Worker's memory.
22const MAX_BYTES: usize = 60 * 1024 * 1024;
23const ARTIFACT_TTL: u64 = 14 * 24 * 60 * 60;
24const CACHE_TTL: u64 = 7 * 24 * 60 * 60;
25
26#[derive(Serialize, Deserialize)]
27struct Meta {
28 size: usize,
29 chunks: usize,
30 /// Milliseconds since the epoch, to find the newest cache entry.
31 at: u64,
32 /// The name or key it was saved under.
33 name: String,
34}
35
36fn store(env: &Env) -> Result<KvStore> {
37 env.kv("BLOBS")
38}
39
40async fn put(kv: &KvStore, base: &str, name: &str, bytes: &[u8], ttl: u64) -> Result<()> {
41 let chunks: Vec<&[u8]> = if bytes.is_empty() { vec![&[][..]] } else { bytes.chunks(CHUNK).collect() };
42 for (index, chunk) in chunks.iter().enumerate() {
43 kv.put_bytes(&format!("{base}#{index}"), chunk)?.expiration_ttl(ttl).execute().await?;
44 }
45 let meta = Meta { size: bytes.len(), chunks: chunks.len(), at: g1t_kit::now_ms(), name: name.to_owned() };
46 // The metadata travels with the key in listings, so the newest entry
47 // can be found without reading each.
48 kv.put(base, serde_json::to_string(&meta)?)?
49 .metadata(&meta)?
50 .expiration_ttl(ttl)
51 .execute()
52 .await?;
53 Ok(())
54}
55
56async fn get(kv: &KvStore, base: &str) -> Result<Option<Vec<u8>>> {
57 let Some(meta) = kv.get(base).json::<Meta>().await? else { return Ok(None) };
58 let mut out = Vec::with_capacity(meta.size);
59 for index in 0..meta.chunks {
60 match kv.get(&format!("{base}#{index}")).bytes().await? {
61 Some(bytes) => out.extend_from_slice(&bytes),
62 // A chunk that expired first: the whole entry is gone.
63 None => return Ok(None),
64 }
65 }
66 Ok(Some(out))
67}
68
69/// The metadata of every entry under a prefix, newest first.
70async fn list(kv: &KvStore, prefix: &str) -> Result<Vec<(String, Meta)>> {
71 let mut found = Vec::new();
72 let mut cursor: Option<String> = None;
73 loop {
74 let mut listing = kv.list().prefix(prefix.to_owned());
75 if let Some(cursor) = cursor.take() {
76 listing = listing.cursor(cursor);
77 }
78 let page = listing.execute().await?;
79 for key in page.keys {
80 if key.name.contains('#') {
81 continue;
82 }
83 if let Some(meta) = key.metadata.and_then(|m| serde_json::from_value::<Meta>(m).ok()) {
84 found.push((key.name, meta));
85 }
86 }
87 if page.list_complete || page.cursor.is_none() {
88 break;
89 }
90 cursor = page.cursor;
91 }
92 found.sort_by_key(|entry| std::cmp::Reverse(entry.1.at));
93 Ok(found)
94}
95
96fn error(status: u16, message: &str) -> Result<Response> {
97 Ok(Response::from_json(&json!({ "error": { "message": message } }))?.with_status(status))
98}
99
100fn valid_name(name: &str) -> bool {
101 !name.is_empty() && name.len() <= 200 && !name.starts_with('.') && name.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.' | ' '))
102}
103
104fn decode(text: &str) -> String {
105 let bytes = text.as_bytes();
106 let mut out = Vec::with_capacity(bytes.len());
107 let mut i = 0;
108 while i < bytes.len() {
109 match bytes[i] {
110 b'%' if i + 2 < bytes.len() => {
111 match u8::from_str_radix(std::str::from_utf8(&bytes[i + 1..i + 3]).unwrap_or("zz"), 16) {
112 Ok(byte) => {
113 out.push(byte);
114 i += 3;
115 }
116 Err(_) => {
117 out.push(b'%');
118 i += 1;
119 }
120 }
121 }
122 b'+' => {
123 out.push(b' ');
124 i += 1;
125 }
126 byte => {
127 out.push(byte);
128 i += 1;
129 }
130 }
131 }
132 String::from_utf8_lossy(&out).into_owned()
133}
134
135fn query(request: &Request, name: &str) -> Option<String> {
136 let url = request.url().ok()?;
137 url.query_pairs().find(|(key, _)| key == name).map(|(_, value)| value.into_owned())
138}
139
140/// A sandbox storing or fetching an artifact or cache entry. `rest` is
141/// the path after `/actions/jobs/`.
142pub async fn for_job(mut request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
143 let (job, what) = rest.split_once('/').unwrap_or((rest, ""));
144 let token = request
145 .headers()
146 .get("authorization")?
147 .and_then(|h| h.strip_prefix("Bearer ").map(str::to_owned))
148 .unwrap_or_default();
149 let owner: Outcome<Value> = g1t_kit::call(&services.actions, "job_auth", &json!({ "job": job, "token": token })).await?;
150 let owner = match owner {
151 Outcome::Ok(owner) => owner,
152 Outcome::Fail(_) => return error(401, "That job is not running, or the token is not its."),
153 };
154 let run = owner["run"].as_str().unwrap_or_default().to_owned();
155 let repo = owner["repoId"].as_str().unwrap_or_default().to_owned();
156 let kv = store(env)?;
157 match (method, what) {
158 ("GET", "artifacts") => {
159 let listed: Vec<Value> = list(&kv, &format!("a/{run}/"))
160 .await?
161 .into_iter()
162 .map(|(_, meta)| json!({ "name": meta.name, "size": meta.size }))
163 .collect();
164 Response::from_json(&listed)
165 }
166 (_, what) if what.starts_with("artifacts/") => {
167 let name = decode(&what["artifacts/".len()..]);
168 if !valid_name(&name) {
169 return error(400, "That is not an artifact name.");
170 }
171 let base = format!("a/{run}/{name}");
172 if method == "PUT" {
173 let bytes = request.bytes().await?;
174 if bytes.len() > MAX_BYTES {
175 return error(413, "Artifacts are at most 60 MB.");
176 }
177 put(&kv, &base, &name, &bytes, ARTIFACT_TTL).await?;
178 return Response::from_json(&json!({ "name": name, "size": bytes.len() }));
179 }
180 match get(&kv, &base).await? {
181 Some(bytes) => Response::from_bytes(bytes),
182 None => error(404, "No such artifact."),
183 }
184 }
185 (_, "cache") => {
186 let key = query(&request, "key").unwrap_or_default();
187 if key.is_empty() || key.len() > 400 {
188 return error(400, "A cache key is 1 to 400 characters.");
189 }
190 if method == "PUT" {
191 let base = format!("c/{repo}/{key}");
192 // A key is written once, as on GitHub.
193 if kv.get(&base).text().await?.is_some() {
194 return Response::from_json(&json!({ "saved": false, "reason": "That key is already cached." }));
195 }
196 let bytes = request.bytes().await?;
197 if bytes.len() > MAX_BYTES {
198 return error(413, "Cache entries are at most 60 MB.");
199 }
200 put(&kv, &base, &key, &bytes, CACHE_TTL).await?;
201 return Response::from_json(&json!({ "saved": true }));
202 }
203 // The exact key, else the newest entry under each restore key.
204 let exact = format!("c/{repo}/{key}");
205 if let Some(bytes) = get(&kv, &exact).await? {
206 let mut response = Response::from_bytes(bytes)?;
207 response.headers_mut().set("x-g1t-key", &key)?;
208 return Ok(response);
209 }
210 for prefix in query(&request, "restore").unwrap_or_default().lines().map(str::trim).filter(|p| !p.is_empty()) {
211 if let Some((base, meta)) = list(&kv, &format!("c/{repo}/{prefix}")).await?.into_iter().next()
212 && let Some(bytes) = get(&kv, &base).await?
213 {
214 let mut response = Response::from_bytes(bytes)?;
215 response.headers_mut().set("x-g1t-key", &meta.name)?;
216 return Ok(response);
217 }
218 }
219 error(404, "Nothing cached under those keys.")
220 }
221 _ => error(404, "No such endpoint."),
222 }
223}
224
225/// Someone who can see the run downloading one of its artifacts.
226pub async fn download(env: &Env, services: &Services, viewer: &g1t_contracts::Viewer, owner: &str, repo: &str, run: &str, name: &str) -> Result<Response> {
227 let seen: Outcome<Value> = g1t_kit::call(
228 &services.actions,
229 "run",
230 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
231 )
232 .await?;
233 if let Outcome::Fail(refused) = seen {
234 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
235 return error(status, &refused.message);
236 }
237 let name = decode(name);
238 match get(&store(env)?, &format!("a/{run}/{name}")).await? {
239 Some(bytes) => {
240 let mut response = Response::from_bytes(bytes)?;
241 let headers = response.headers_mut();
242 headers.set("content-type", "application/gzip")?;
243 headers.set("content-disposition", &format!("attachment; filename=\"{}.tar.gz\"", name.replace('"', "")))?;
244 Ok(response)
245 }
246 None => error(404, "No such artifact, or it has expired."),
247 }
248}
249
250/// A run's artifacts, for its page.
251pub async fn of_run(env: &Env, run: &str) -> Result<Vec<Value>> {
252 Ok(list(&store(env)?, &format!("a/{run}/"))
253 .await?
254 .into_iter()
255 .map(|(_, meta)| json!({ "name": meta.name, "size": meta.size, "at": meta.at }))
256 .collect())
257}