Skip to content
434 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Fast pages, required checks on the branch, self-hosted runners, honest incidents1//! 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)
A repository has its own sidebar, as settings do19//!
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;
Fast pages, required checks on the branch, self-hosted runners, honest incidents26use worker::{Bucket, Env, Request, Response, Result, UploadedPart};
A repository has its own sidebar, as settings do27
Fast pages, required checks on the branch, self-hosted runners, honest incidents28use g1t_contracts::actions::{
29 CACHE_PART_BYTES, CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs,
30};
A repository has its own sidebar, as settings do31use 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> {
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API112 Ok(crate::reply(&json!({ "error": { "message": message } }))?.with_status(status))
A repository has its own sidebar, as settings do113}
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();
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API179 crate::reply(&listed)
A repository has its own sidebar, as settings do180 }
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?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API193 return crate::reply(&json!({ "name": name, "size": bytes.len() }));
A repository has its own sidebar, as settings do194 }
195 match get(&kv, &base).await? {
196 Some(bytes) => Response::from_bytes(bytes),
197 None => error(404, "No such artifact."),
198 }
199 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents200 (_, what) if what == "cache" || what.starts_with("cache/") => {
201 let bucket = env.bucket("ACTIONS_CACHE")?;
Actions: keep workflow runs safe202 cache(request, &bucket, services, method, job, &token, &repo, what).await
Fast pages, required checks on the branch, self-hosted runners, honest incidents203 }
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 bucket: &Bucket,
243 services: &Services,
244 method: &str,
245 job: &str,
246 token: &str,
247 repo: &str,
248 what: &str,
249) -> Result<Response> {
250 let parts: Vec<&str> = what.split('/').collect();
251 let upload_id = query(&request, "upload").unwrap_or_default();
Actions: keep workflow runs safe252 // The hash of the entry's paths and compression; runners from before
253 // it was sent send none.
254 let version = query(&request, "version").unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents255 match (method, parts.as_slice()) {
256 ("GET", ["cache"]) => {
257 let key = query(&request, "key").unwrap_or_default();
258 let restore: Vec<String> =
259 query(&request, "restore").unwrap_or_default().lines().map(str::trim).filter(|p| !p.is_empty()).map(str::to_owned).collect();
260 let found: Outcome<Option<CacheHit>> = g1t_kit::call(
261 &services.actions,
262 "cache_lookup",
Actions: keep workflow runs safe263 &CacheLookupArgs { job: job.to_owned(), token: token.to_owned(), key: key.clone(), restore, version: version.clone() },
Fast pages, required checks on the branch, self-hosted runners, honest incidents264 )
265 .await?;
266 let found = match refused(found) {
267 Ok(found) => found,
268 Err(reply) => return reply,
269 };
270 if let Some(hit) = found
271 && let Some(object) = bucket.get(&hit.object).execute().await?
272 && let Some(body) = object.body()
273 {
274 let mut response = Response::from_body(body.response_body()?)?;
275 let headers = response.headers_mut();
276 headers.set("x-g1t-key", &hit.key)?;
277 headers.set("content-length", &object.size().to_string())?;
278 headers.set("content-type", "application/octet-stream")?;
279 return Ok(response);
280 }
Actions: keep workflow runs safe281 // Entries kept in KV before the cache moved to R2 had no scope,
282 // so they are never restored.
283 error(404, "Nothing cached under those keys.")
Fast pages, required checks on the branch, self-hosted runners, honest incidents284 }
285 // Older runners send a whole entry of at most 60 MB at once.
286 ("PUT", ["cache"]) => {
A repository has its own sidebar, as settings do287 let key = query(&request, "key").unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents288 let bytes = request.bytes().await?;
289 if bytes.len() > MAX_BYTES {
290 return error(413, "An entry sent at once is at most 60 MB; newer runners upload it in parts.");
A repository has its own sidebar, as settings do291 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents292 let reserved: Outcome<CacheReservation> = g1t_kit::call(
293 &services.actions,
294 "cache_reserve",
Actions: keep workflow runs safe295 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size: bytes.len() as u64, version: version.clone() },
Fast pages, required checks on the branch, self-hosted runners, honest incidents296 )
297 .await?;
298 let reserved = match reserved {
299 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
300 return crate::reply(&json!({ "saved": false, "reason": failure.message }));
A repository has its own sidebar, as settings do301 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents302 other => match refused(other) {
303 Ok(reserved) => reserved,
304 Err(reply) => return reply,
305 },
306 };
307 let size = bytes.len() as u64;
308 bucket.put(&reserved.object, bytes).execute().await?;
309 commit(bucket, services, job, token, &reserved.id, size).await?;
310 crate::reply(&json!({ "saved": true }))
311 }
312 ("POST", ["cache", "uploads"]) => {
313 let key = query(&request, "key").unwrap_or_default();
314 let size = query(&request, "size").and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
315 let reserved: Outcome<CacheReservation> = g1t_kit::call(
316 &services.actions,
317 "cache_reserve",
Actions: keep workflow runs safe318 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size, version: version.clone() },
Fast pages, required checks on the branch, self-hosted runners, honest incidents319 )
320 .await?;
321 let reserved = match refused(reserved) {
322 Ok(reserved) => reserved,
323 Err(reply) => return reply,
324 };
325 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
326 crate::reply(&json!({ "id": reserved.id, "upload": upload.upload_id().await, "part_bytes": CACHE_PART_BYTES }))
327 }
328 ("PUT", ["cache", "uploads", id, part]) => {
329 let part = part.parse::<u16>().unwrap_or(0);
330 if part == 0 || upload_id.is_empty() {
331 return error(400, "A part is numbered from 1, and names its upload.");
A repository has its own sidebar, as settings do332 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents333 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
334 if length == 0 || length > CACHE_PART_BYTES {
335 return error(413, &format!("A part is 1 to {} MB, with its length.", CACHE_PART_BYTES / 1_048_576));
A repository has its own sidebar, as settings do336 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents337 // The body goes to R2 as it comes, never held whole here.
338 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
339 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
340 let uploaded = upload.upload_part(part, body).await?;
341 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
342 }
343 ("POST", ["cache", "uploads", id, "complete"]) => {
344 let done: Complete = match request.json().await {
345 Ok(done) => done,
346 Err(_) => return error(400, "Send { size, parts: [{ part, etag }] }."),
347 };
348 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
349 let mut parts = done.parts;
350 parts.sort_by_key(|p| p.part);
351 if let Err(problem) = upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
352 let _ = abort(services, job, token, id).await;
353 return error(400, &format!("The upload could not be completed: {problem}"));
354 }
355 commit(bucket, services, job, token, id, done.size).await?;
356 crate::reply(&json!({ "saved": true }))
357 }
358 ("DELETE", ["cache", "uploads", id]) => {
359 if let Ok(upload) = bucket.resume_multipart_upload(object_of(repo, id), &upload_id) {
360 let _ = upload.abort().await;
A repository has its own sidebar, as settings do361 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents362 abort(services, job, token, id).await?;
363 crate::reply(&json!({ "aborted": true }))
A repository has its own sidebar, as settings do364 }
365 _ => error(404, "No such endpoint."),
366 }
367}
368
Fast pages, required checks on the branch, self-hosted runners, honest incidents369/// Where an entry is in R2: under its repository, by its id, as the
370/// actions service named it when it was reserved.
371fn object_of(repo: &str, id: &str) -> String {
372 format!("c/{repo}/{id}")
373}
374
375/// Marks an uploaded entry ready, and deletes what that evicted.
376async fn commit(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<()> {
377 let committed: Outcome<CacheCommitted> = g1t_kit::call(
378 &services.actions,
379 "cache_commit",
380 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
381 )
382 .await?;
383 if let Outcome::Ok(committed) = committed
384 && !committed.evicted.is_empty()
385 {
386 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
387 }
388 Ok(())
389}
390
391async fn abort(services: &Services, job: &str, token: &str, id: &str) -> Result<()> {
392 let _: Outcome<bool> = g1t_kit::call(
393 &services.actions,
394 "cache_abort",
395 &CacheAbortArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned() },
396 )
397 .await?;
398 Ok(())
399}
400
401/// An entry saved in KV before the cache moved to R2, by key or restore key.
A repository has its own sidebar, as settings do402/// Someone who can see the run downloading one of its artifacts.
403pub async fn download(env: &Env, services: &Services, viewer: &g1t_contracts::Viewer, owner: &str, repo: &str, run: &str, name: &str) -> Result<Response> {
404 let seen: Outcome<Value> = g1t_kit::call(
405 &services.actions,
406 "run",
407 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
408 )
409 .await?;
410 if let Outcome::Fail(refused) = seen {
411 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
412 return error(status, &refused.message);
413 }
414 let name = decode(name);
415 match get(&store(env)?, &format!("a/{run}/{name}")).await? {
416 Some(bytes) => {
417 let mut response = Response::from_bytes(bytes)?;
418 let headers = response.headers_mut();
419 headers.set("content-type", "application/gzip")?;
420 headers.set("content-disposition", &format!("attachment; filename=\"{}.tar.gz\"", name.replace('"', "")))?;
421 Ok(response)
422 }
423 None => error(404, "No such artifact, or it has expired."),
424 }
425}
426
427/// A run's artifacts, for its page.
428pub async fn of_run(env: &Env, run: &str) -> Result<Vec<Value>> {
429 Ok(list(&store(env)?, &format!("a/{run}/"))
430 .await?
431 .into_iter()
432 .map(|(_, meta)| json!({ "name": meta.name, "size": meta.size, "at": meta.at }))
433 .collect())
434}