flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/apps/api/src/blobs.rs

450 lines18,637 bytesCodeBlame
1//! Artifacts and the cache of GitHub Actions jobs.
2//!
3//! Artifacts are kept in Workers KV in chunks, with KV's own expiry: 14
4//! days with their run. Cache entries are kept in R2 (ACTIONS_CACHE), up
5//! to 2 GB each, uploaded in parts; the actions service lists them and
6//! decides what is found, what fits and what is evicted
7//! (services/actions/src/cache.rs). Entries saved in KV before the cache
8//! moved are still found there until they expire.
9//!
10//! A sandbox reaches these with its job's token:
11//!
12//! - `GET /actions/jobs/{job}/artifacts`, `PUT|GET .../artifacts/{name}`
13//! - `GET .../cache?key=&restore=`: the entry, streamed, its key in `x-g1t-key`
14//! - `POST .../cache/uploads?key=&size=`: `{ id, upload, part_bytes }`
15//! - `PUT .../cache/uploads/{id}/{part}?upload=`: one part, `{ part, etag }`
16//! - `POST .../cache/uploads/{id}/complete?upload=` with `{ size, parts }`
17//! - `DELETE .../cache/uploads/{id}?upload=`: gives the upload up
18//! - `PUT .../cache?key=`: a whole entry of at most 60 MB at once (older runners)
19//!
20//! People download an artifact at
21//! `/repos/{owner}/{repo}/actions/runs/{run}/artifacts/{name}`.
22
23use serde::{Deserialize, Serialize};
24use serde_json::{Value, json};
25use worker::kv::KvStore;
26use worker::{Bucket, Env, Request, Response, Result, UploadedPart};
27
28use g1t_contracts::actions::{
29 CACHE_PART_BYTES, CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs,
30};
31use g1t_contracts::{FailureCode, Outcome};
32
33use crate::operations::Services;
34
35/// KV's largest value is 25 MiB; chunks stay under it.
36const CHUNK: usize = 20 * 1024 * 1024;
37/// The largest artifact or cache entry, kept within a Worker's memory.
38const MAX_BYTES: usize = 60 * 1024 * 1024;
39const ARTIFACT_TTL: u64 = 14 * 24 * 60 * 60;
40
41#[derive(Serialize, Deserialize)]
42struct Meta {
43 size: usize,
44 chunks: usize,
45 /// Milliseconds since the epoch, to find the newest cache entry.
46 at: u64,
47 /// The name or key it was saved under.
48 name: String,
49}
50
51fn store(env: &Env) -> Result<KvStore> {
52 env.kv("BLOBS")
53}
54
55async fn put(kv: &KvStore, base: &str, name: &str, bytes: &[u8], ttl: u64) -> Result<()> {
56 let chunks: Vec<&[u8]> = if bytes.is_empty() { vec![&[][..]] } else { bytes.chunks(CHUNK).collect() };
57 for (index, chunk) in chunks.iter().enumerate() {
58 kv.put_bytes(&format!("{base}#{index}"), chunk)?.expiration_ttl(ttl).execute().await?;
59 }
60 let meta = Meta { size: bytes.len(), chunks: chunks.len(), at: g1t_kit::now_ms(), name: name.to_owned() };
61 // The metadata travels with the key in listings, so the newest entry
62 // can be found without reading each.
63 kv.put(base, serde_json::to_string(&meta)?)?
64 .metadata(&meta)?
65 .expiration_ttl(ttl)
66 .execute()
67 .await?;
68 Ok(())
69}
70
71async fn get(kv: &KvStore, base: &str) -> Result<Option<Vec<u8>>> {
72 let Some(meta) = kv.get(base).json::<Meta>().await? else { return Ok(None) };
73 let mut out = Vec::with_capacity(meta.size);
74 for index in 0..meta.chunks {
75 match kv.get(&format!("{base}#{index}")).bytes().await? {
76 Some(bytes) => out.extend_from_slice(&bytes),
77 // A chunk that expired first: the whole entry is gone.
78 None => return Ok(None),
79 }
80 }
81 Ok(Some(out))
82}
83
84/// The metadata of every entry under a prefix, newest first.
85async fn list(kv: &KvStore, prefix: &str) -> Result<Vec<(String, Meta)>> {
86 let mut found = Vec::new();
87 let mut cursor: Option<String> = None;
88 loop {
89 let mut listing = kv.list().prefix(prefix.to_owned());
90 if let Some(cursor) = cursor.take() {
91 listing = listing.cursor(cursor);
92 }
93 let page = listing.execute().await?;
94 for key in page.keys {
95 if key.name.contains('#') {
96 continue;
97 }
98 if let Some(meta) = key.metadata.and_then(|m| serde_json::from_value::<Meta>(m).ok()) {
99 found.push((key.name, meta));
100 }
101 }
102 if page.list_complete || page.cursor.is_none() {
103 break;
104 }
105 cursor = page.cursor;
106 }
107 found.sort_by_key(|entry| std::cmp::Reverse(entry.1.at));
108 Ok(found)
109}
110
111fn error(status: u16, message: &str) -> Result<Response> {
112 Ok(crate::reply(&json!({ "error": { "message": message } }))?.with_status(status))
113}
114
115fn valid_name(name: &str) -> bool {
116 !name.is_empty() && name.len() <= 200 && !name.starts_with('.') && name.chars().all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.' | ' '))
117}
118
119fn decode(text: &str) -> String {
120 let bytes = text.as_bytes();
121 let mut out = Vec::with_capacity(bytes.len());
122 let mut i = 0;
123 while i < bytes.len() {
124 match bytes[i] {
125 b'%' if i + 2 < bytes.len() => {
126 match u8::from_str_radix(std::str::from_utf8(&bytes[i + 1..i + 3]).unwrap_or("zz"), 16) {
127 Ok(byte) => {
128 out.push(byte);
129 i += 3;
130 }
131 Err(_) => {
132 out.push(b'%');
133 i += 1;
134 }
135 }
136 }
137 b'+' => {
138 out.push(b' ');
139 i += 1;
140 }
141 byte => {
142 out.push(byte);
143 i += 1;
144 }
145 }
146 }
147 String::from_utf8_lossy(&out).into_owned()
148}
149
150fn query(request: &Request, name: &str) -> Option<String> {
151 let url = request.url().ok()?;
152 url.query_pairs().find(|(key, _)| key == name).map(|(_, value)| value.into_owned())
153}
154
155/// A sandbox storing or fetching an artifact or cache entry. `rest` is
156/// the path after `/actions/jobs/`.
157pub async fn for_job(mut request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
158 let (job, what) = rest.split_once('/').unwrap_or((rest, ""));
159 let token = request
160 .headers()
161 .get("authorization")?
162 .and_then(|h| h.strip_prefix("Bearer ").map(str::to_owned))
163 .unwrap_or_default();
164 let owner: Outcome<Value> = g1t_kit::call(&services.actions, "job_auth", &json!({ "job": job, "token": token })).await?;
165 let owner = match owner {
166 Outcome::Ok(owner) => owner,
167 Outcome::Fail(_) => return error(401, "That job is not running, or the token is not its."),
168 };
169 let run = owner["run"].as_str().unwrap_or_default().to_owned();
170 let repo = owner["repoId"].as_str().unwrap_or_default().to_owned();
171 let kv = store(env)?;
172 match (method, what) {
173 ("GET", "artifacts") => {
174 let listed: Vec<Value> = list(&kv, &format!("a/{run}/"))
175 .await?
176 .into_iter()
177 .map(|(_, meta)| json!({ "name": meta.name, "size": meta.size }))
178 .collect();
179 crate::reply(&listed)
180 }
181 (_, what) if what.starts_with("artifacts/") => {
182 let name = decode(&what["artifacts/".len()..]);
183 if !valid_name(&name) {
184 return error(400, "That is not an artifact name.");
185 }
186 let base = format!("a/{run}/{name}");
187 if method == "PUT" {
188 let bytes = request.bytes().await?;
189 if bytes.len() > MAX_BYTES {
190 return error(413, "Artifacts are at most 60 MB.");
191 }
192 put(&kv, &base, &name, &bytes, ARTIFACT_TTL).await?;
193 return crate::reply(&json!({ "name": name, "size": bytes.len() }));
194 }
195 match get(&kv, &base).await? {
196 Some(bytes) => Response::from_bytes(bytes),
197 None => error(404, "No such artifact."),
198 }
199 }
200 (_, what) if what == "cache" || what.starts_with("cache/") => {
201 let bucket = env.bucket("ACTIONS_CACHE")?;
202 cache(request, &kv, &bucket, services, method, job, &token, &repo, what).await
203 }
204 _ => error(404, "No such endpoint."),
205 }
206}
207
208/// A part's number and the etag R2 gave it.
209#[derive(Deserialize)]
210struct Part {
211 part: u16,
212 etag: String,
213}
214
215#[derive(Deserialize)]
216struct Complete {
217 size: u64,
218 parts: Vec<Part>,
219}
220
221/// An outcome of the actions service, or its failure as the reply it means.
222fn refused<T>(outcome: Outcome<T>) -> std::result::Result<T, Result<Response>> {
223 match outcome {
224 Outcome::Ok(value) => Ok(value),
225 Outcome::Fail(failure) => {
226 let status = match failure.code {
227 FailureCode::Unauthenticated => 401,
228 FailureCode::NotFound => 404,
229 FailureCode::Conflict => 409,
230 FailureCode::Forbidden => 403,
231 _ => 400,
232 };
233 Err(error(status, &failure.message))
234 }
235 }
236}
237
238/// The cache: restoring, uploading in parts, and the older whole upload.
239#[allow(clippy::too_many_arguments)]
240async fn cache(
241 mut request: Request,
242 kv: &KvStore,
243 bucket: &Bucket,
244 services: &Services,
245 method: &str,
246 job: &str,
247 token: &str,
248 repo: &str,
249 what: &str,
250) -> Result<Response> {
251 let parts: Vec<&str> = what.split('/').collect();
252 let upload_id = query(&request, "upload").unwrap_or_default();
253 match (method, parts.as_slice()) {
254 ("GET", ["cache"]) => {
255 let key = query(&request, "key").unwrap_or_default();
256 let restore: Vec<String> =
257 query(&request, "restore").unwrap_or_default().lines().map(str::trim).filter(|p| !p.is_empty()).map(str::to_owned).collect();
258 let found: Outcome<Option<CacheHit>> = g1t_kit::call(
259 &services.actions,
260 "cache_lookup",
261 &CacheLookupArgs { job: job.to_owned(), token: token.to_owned(), key: key.clone(), restore: restore.clone() },
262 )
263 .await?;
264 let found = match refused(found) {
265 Ok(found) => found,
266 Err(reply) => return reply,
267 };
268 if let Some(hit) = found
269 && let Some(object) = bucket.get(&hit.object).execute().await?
270 && let Some(body) = object.body()
271 {
272 let mut response = Response::from_body(body.response_body()?)?;
273 let headers = response.headers_mut();
274 headers.set("x-g1t-key", &hit.key)?;
275 headers.set("content-length", &object.size().to_string())?;
276 headers.set("content-type", "application/octet-stream")?;
277 return Ok(response);
278 }
279 // Entries saved in KV before the cache moved to R2.
280 kv_lookup(kv, repo, &key, &restore).await
281 }
282 // Older runners send a whole entry of at most 60 MB at once.
283 ("PUT", ["cache"]) => {
284 let key = query(&request, "key").unwrap_or_default();
285 let bytes = request.bytes().await?;
286 if bytes.len() > MAX_BYTES {
287 return error(413, "An entry sent at once is at most 60 MB; newer runners upload it in parts.");
288 }
289 let reserved: Outcome<CacheReservation> = g1t_kit::call(
290 &services.actions,
291 "cache_reserve",
292 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size: bytes.len() as u64 },
293 )
294 .await?;
295 let reserved = match reserved {
296 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
297 return crate::reply(&json!({ "saved": false, "reason": failure.message }));
298 }
299 other => match refused(other) {
300 Ok(reserved) => reserved,
301 Err(reply) => return reply,
302 },
303 };
304 let size = bytes.len() as u64;
305 bucket.put(&reserved.object, bytes).execute().await?;
306 commit(bucket, services, job, token, &reserved.id, size).await?;
307 crate::reply(&json!({ "saved": true }))
308 }
309 ("POST", ["cache", "uploads"]) => {
310 let key = query(&request, "key").unwrap_or_default();
311 let size = query(&request, "size").and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
312 let reserved: Outcome<CacheReservation> = g1t_kit::call(
313 &services.actions,
314 "cache_reserve",
315 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size },
316 )
317 .await?;
318 let reserved = match refused(reserved) {
319 Ok(reserved) => reserved,
320 Err(reply) => return reply,
321 };
322 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
323 crate::reply(&json!({ "id": reserved.id, "upload": upload.upload_id().await, "part_bytes": CACHE_PART_BYTES }))
324 }
325 ("PUT", ["cache", "uploads", id, part]) => {
326 let part = part.parse::<u16>().unwrap_or(0);
327 if part == 0 || upload_id.is_empty() {
328 return error(400, "A part is numbered from 1, and names its upload.");
329 }
330 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
331 if length == 0 || length > CACHE_PART_BYTES {
332 return error(413, &format!("A part is 1 to {} MB, with its length.", CACHE_PART_BYTES / 1_048_576));
333 }
334 // The body goes to R2 as it comes, never held whole here.
335 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
336 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
337 let uploaded = upload.upload_part(part, body).await?;
338 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
339 }
340 ("POST", ["cache", "uploads", id, "complete"]) => {
341 let done: Complete = match request.json().await {
342 Ok(done) => done,
343 Err(_) => return error(400, "Send { size, parts: [{ part, etag }] }."),
344 };
345 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
346 let mut parts = done.parts;
347 parts.sort_by_key(|p| p.part);
348 if let Err(problem) = upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
349 let _ = abort(services, job, token, id).await;
350 return error(400, &format!("The upload could not be completed: {problem}"));
351 }
352 commit(bucket, services, job, token, id, done.size).await?;
353 crate::reply(&json!({ "saved": true }))
354 }
355 ("DELETE", ["cache", "uploads", id]) => {
356 if let Ok(upload) = bucket.resume_multipart_upload(object_of(repo, id), &upload_id) {
357 let _ = upload.abort().await;
358 }
359 abort(services, job, token, id).await?;
360 crate::reply(&json!({ "aborted": true }))
361 }
362 _ => error(404, "No such endpoint."),
363 }
364}
365
366/// Where an entry is in R2: under its repository, by its id, as the
367/// actions service named it when it was reserved.
368fn object_of(repo: &str, id: &str) -> String {
369 format!("c/{repo}/{id}")
370}
371
372/// Marks an uploaded entry ready, and deletes what that evicted.
373async fn commit(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<()> {
374 let committed: Outcome<CacheCommitted> = g1t_kit::call(
375 &services.actions,
376 "cache_commit",
377 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
378 )
379 .await?;
380 if let Outcome::Ok(committed) = committed
381 && !committed.evicted.is_empty()
382 {
383 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
384 }
385 Ok(())
386}
387
388async fn abort(services: &Services, job: &str, token: &str, id: &str) -> Result<()> {
389 let _: Outcome<bool> = g1t_kit::call(
390 &services.actions,
391 "cache_abort",
392 &CacheAbortArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned() },
393 )
394 .await?;
395 Ok(())
396}
397
398/// An entry saved in KV before the cache moved to R2, by key or restore key.
399async fn kv_lookup(kv: &KvStore, repo: &str, key: &str, restore: &[String]) -> Result<Response> {
400 // The exact key, else the newest entry under each restore key.
401 if let Some(bytes) = get(kv, &format!("c/{repo}/{key}")).await? {
402 let mut response = Response::from_bytes(bytes)?;
403 response.headers_mut().set("x-g1t-key", key)?;
404 return Ok(response);
405 }
406 for prefix in restore {
407 if let Some((base, meta)) = list(kv, &format!("c/{repo}/{prefix}")).await?.into_iter().next()
408 && let Some(bytes) = get(kv, &base).await?
409 {
410 let mut response = Response::from_bytes(bytes)?;
411 response.headers_mut().set("x-g1t-key", &meta.name)?;
412 return Ok(response);
413 }
414 }
415 error(404, "Nothing cached under those keys.")
416}
417
418/// Someone who can see the run downloading one of its artifacts.
419pub async fn download(env: &Env, services: &Services, viewer: &g1t_contracts::Viewer, owner: &str, repo: &str, run: &str, name: &str) -> Result<Response> {
420 let seen: Outcome<Value> = g1t_kit::call(
421 &services.actions,
422 "run",
423 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
424 )
425 .await?;
426 if let Outcome::Fail(refused) = seen {
427 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
428 return error(status, &refused.message);
429 }
430 let name = decode(name);
431 match get(&store(env)?, &format!("a/{run}/{name}")).await? {
432 Some(bytes) => {
433 let mut response = Response::from_bytes(bytes)?;
434 let headers = response.headers_mut();
435 headers.set("content-type", "application/gzip")?;
436 headers.set("content-disposition", &format!("attachment; filename=\"{}.tar.gz\"", name.replace('"', "")))?;
437 Ok(response)
438 }
439 None => error(404, "No such artifact, or it has expired."),
440 }
441}
442
443/// A run's artifacts, for its page.
444pub async fn of_run(env: &Env, run: &str) -> Result<Vec<Value>> {
445 Ok(list(&store(env)?, &format!("a/{run}/"))
446 .await?
447 .into_iter()
448 .map(|(_, meta)| json!({ "name": meta.name, "size": meta.size, "at": meta.at }))
449 .collect())
450}