pr_01m47d15m3e54sn21z27rpy5n9/apps/api/src/blobs.rs
| 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 | |
| 10 | use serde::{Deserialize, Serialize}; |
| 11 | use serde_json::{Value, json}; |
| 12 | use worker::kv::KvStore; |
| 13 | use worker::{Env, Request, Response, Result}; |
| 14 | |
| 15 | use g1t_contracts::{FailureCode, Outcome}; |
| 16 | |
| 17 | use crate::operations::Services; |
| 18 | |
| 19 | /// KV's largest value is 25 MiB; chunks stay under it. |
| 20 | const CHUNK: usize = 20 * 1024 * 1024; |
| 21 | /// The largest artifact or cache entry, kept within a Worker's memory. |
| 22 | const MAX_BYTES: usize = 60 * 1024 * 1024; |
| 23 | const ARTIFACT_TTL: u64 = 14 * 24 * 60 * 60; |
| 24 | const CACHE_TTL: u64 = 7 * 24 * 60 * 60; |
| 25 | |
| 26 | #[derive(Serialize, Deserialize)] |
| 27 | struct 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 | |
| 36 | fn store(env: &Env) -> Result<KvStore> { |
| 37 | env.kv("BLOBS") |
| 38 | } |
| 39 | |
| 40 | async 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 | |
| 56 | async 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. |
| 70 | async 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 | |
| 96 | fn error(status: u16, message: &str) -> Result<Response> { |
| 97 | Ok(Response::from_json(&json!({ "error": { "message": message } }))?.with_status(status)) |
| 98 | } |
| 99 | |
| 100 | fn 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 | |
| 104 | fn 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 | |
| 135 | fn 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/`. |
| 142 | pub 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. |
| 226 | pub 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. |
| 251 | pub 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 | } |