Skip to content
626 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.

Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R21//! Artifacts and the cache of GitHub Actions jobs, as g1t's own runner
2//! reaches them. (Actions built on GitHub's toolkit reach the same entries
3//! through toolkit.rs.)
Fast pages, required checks on the branch, self-hosted runners, honest incidents4//!
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R25//! Artifacts are kept in R2 (ACTIONS_CACHE, under `a/`), uploaded in
6//! parts; the actions service lists them and decides names, sizes and how
7//! long each is kept (services/actions/src/artifacts.rs). Artifacts older
8//! runners kept in Workers KV are still listed and found there until KV
9//! expires them. Cache entries are kept in R2 too, up to 2 GB each,
10//! uploaded in parts; the actions service decides what is found, what
11//! fits and what is evicted (services/actions/src/cache.rs).
Fast pages, required checks on the branch, self-hosted runners, honest incidents12//!
13//! A sandbox reaches these with its job's token:
14//!
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R215//! - `GET /actions/jobs/{job}/artifacts[?run_id=]`: its run's (or another run's)
16//! - `POST .../artifacts/uploads?name=&size=&retention_days=&overwrite=&format=`:
17//! `{ id, upload, part_bytes, retention_days, expires_at }`
18//! - `PUT .../artifacts/uploads/{id}/{part}?upload=`, `POST …/complete` with
19//! `{ size, parts, digest }`, `DELETE .../artifacts/uploads/{id}?upload=`
20//! - `GET .../artifacts/{id}/download[?run_id=]`, `DELETE .../artifacts/{id}`
21//! - `PUT|GET .../artifacts/{name}`: a whole artifact by name (older runners)
Fast pages, required checks on the branch, self-hosted runners, honest incidents22//! - `GET .../cache?key=&restore=`: the entry, streamed, its key in `x-g1t-key`
23//! - `POST .../cache/uploads?key=&size=`: `{ id, upload, part_bytes }`
24//! - `PUT .../cache/uploads/{id}/{part}?upload=`: one part, `{ part, etag }`
25//! - `POST .../cache/uploads/{id}/complete?upload=` with `{ size, parts }`
26//! - `DELETE .../cache/uploads/{id}?upload=`: gives the upload up
27//! - `PUT .../cache?key=`: a whole entry of at most 60 MB at once (older runners)
A repository has its own sidebar, as settings do28//!
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R229//! People download an artifact through the REST API (artifacts.rs), or by
30//! name at `/repos/{owner}/{repo}/actions/runs/{run}/artifacts/{name}`.
A repository has its own sidebar, as settings do31
32use serde::{Deserialize, Serialize};
33use serde_json::{Value, json};
34use worker::kv::KvStore;
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R235use worker::{Bucket, Env, Request, Response, Result, UploadedPart, Url};
A repository has its own sidebar, as settings do36
Fast pages, required checks on the branch, self-hosted runners, honest incidents37use g1t_contracts::actions::{
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R238 ARTIFACT_PART_BYTES, Artifact, ArtifactArgs, ArtifactBlob, ArtifactCommitArgs, ArtifactReservation, ArtifactReserveArgs, CACHE_PART_BYTES,
39 CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs, JobArtifactsArgs,
Fast pages, required checks on the branch, self-hosted runners, honest incidents40};
A repository has its own sidebar, as settings do41use g1t_contracts::{FailureCode, Outcome};
42
43use crate::operations::Services;
44
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R245/// The largest artifact or cache entry an older runner sends at once,
46/// held in a Worker's memory.
A repository has its own sidebar, as settings do47const MAX_BYTES: usize = 60 * 1024 * 1024;
48
49#[derive(Serialize, Deserialize)]
50struct Meta {
51 size: usize,
52 chunks: usize,
53 /// Milliseconds since the epoch, to find the newest cache entry.
54 at: u64,
55 /// The name or key it was saved under.
56 name: String,
57}
58
59fn store(env: &Env) -> Result<KvStore> {
60 env.kv("BLOBS")
61}
62
63async fn get(kv: &KvStore, base: &str) -> Result<Option<Vec<u8>>> {
64 let Some(meta) = kv.get(base).json::<Meta>().await? else { return Ok(None) };
65 let mut out = Vec::with_capacity(meta.size);
66 for index in 0..meta.chunks {
67 match kv.get(&format!("{base}#{index}")).bytes().await? {
68 Some(bytes) => out.extend_from_slice(&bytes),
69 // A chunk that expired first: the whole entry is gone.
70 None => return Ok(None),
71 }
72 }
73 Ok(Some(out))
74}
75
76/// The metadata of every entry under a prefix, newest first.
77async fn list(kv: &KvStore, prefix: &str) -> Result<Vec<(String, Meta)>> {
78 let mut found = Vec::new();
79 let mut cursor: Option<String> = None;
80 loop {
81 let mut listing = kv.list().prefix(prefix.to_owned());
82 if let Some(cursor) = cursor.take() {
83 listing = listing.cursor(cursor);
84 }
85 let page = listing.execute().await?;
86 for key in page.keys {
87 if key.name.contains('#') {
88 continue;
89 }
90 if let Some(meta) = key.metadata.and_then(|m| serde_json::from_value::<Meta>(m).ok()) {
91 found.push((key.name, meta));
92 }
93 }
94 if page.list_complete || page.cursor.is_none() {
95 break;
96 }
97 cursor = page.cursor;
98 }
99 found.sort_by_key(|entry| std::cmp::Reverse(entry.1.at));
100 Ok(found)
101}
102
103fn 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 API104 Ok(crate::reply(&json!({ "error": { "message": message } }))?.with_status(status))
A repository has its own sidebar, as settings do105}
106
107fn decode(text: &str) -> String {
108 let bytes = text.as_bytes();
109 let mut out = Vec::with_capacity(bytes.len());
110 let mut i = 0;
111 while i < bytes.len() {
112 match bytes[i] {
113 b'%' if i + 2 < bytes.len() => {
114 match u8::from_str_radix(std::str::from_utf8(&bytes[i + 1..i + 3]).unwrap_or("zz"), 16) {
115 Ok(byte) => {
116 out.push(byte);
117 i += 3;
118 }
119 Err(_) => {
120 out.push(b'%');
121 i += 1;
122 }
123 }
124 }
125 b'+' => {
126 out.push(b' ');
127 i += 1;
128 }
129 byte => {
130 out.push(byte);
131 i += 1;
132 }
133 }
134 }
135 String::from_utf8_lossy(&out).into_owned()
136}
137
138fn query(request: &Request, name: &str) -> Option<String> {
139 let url = request.url().ok()?;
140 url.query_pairs().find(|(key, _)| key == name).map(|(_, value)| value.into_owned())
141}
142
143/// A sandbox storing or fetching an artifact or cache entry. `rest` is
144/// the path after `/actions/jobs/`.
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2145pub async fn for_job(request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
A repository has its own sidebar, as settings do146 let (job, what) = rest.split_once('/').unwrap_or((rest, ""));
147 let token = request
148 .headers()
149 .get("authorization")?
150 .and_then(|h| h.strip_prefix("Bearer ").map(str::to_owned))
151 .unwrap_or_default();
152 let owner: Outcome<Value> = g1t_kit::call(&services.actions, "job_auth", &json!({ "job": job, "token": token })).await?;
153 let owner = match owner {
154 Outcome::Ok(owner) => owner,
155 Outcome::Fail(_) => return error(401, "That job is not running, or the token is not its."),
156 };
157 let run = owner["run"].as_str().unwrap_or_default().to_owned();
158 let repo = owner["repoId"].as_str().unwrap_or_default().to_owned();
159 let kv = store(env)?;
160 match (method, what) {
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2161 (_, what) if what == "artifacts" || what.starts_with("artifacts/") => {
162 let bucket = env.bucket("ACTIONS_CACHE")?;
163 artifacts(request, &kv, &bucket, services, method, job, &token, &run, &repo, what).await
A repository has its own sidebar, as settings do164 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents165 (_, what) if what == "cache" || what.starts_with("cache/") => {
166 let bucket = env.bucket("ACTIONS_CACHE")?;
Merge branch 'worktree-agent-a3abfcce648e87dca'167 cache(request, &bucket, services, method, job, &token, &repo, what).await
Fast pages, required checks on the branch, self-hosted runners, honest incidents168 }
169 _ => error(404, "No such endpoint."),
170 }
171}
172
173/// A part's number and the etag R2 gave it.
174#[derive(Deserialize)]
175struct Part {
176 part: u16,
177 etag: String,
178}
179
180#[derive(Deserialize)]
181struct Complete {
182 size: u64,
183 parts: Vec<Part>,
184}
185
186/// An outcome of the actions service, or its failure as the reply it means.
187fn refused<T>(outcome: Outcome<T>) -> std::result::Result<T, Result<Response>> {
188 match outcome {
189 Outcome::Ok(value) => Ok(value),
190 Outcome::Fail(failure) => {
191 let status = match failure.code {
192 FailureCode::Unauthenticated => 401,
193 FailureCode::NotFound => 404,
194 FailureCode::Conflict => 409,
195 FailureCode::Forbidden => 403,
196 _ => 400,
197 };
198 Err(error(status, &failure.message))
199 }
200 }
201}
202
203/// The cache: restoring, uploading in parts, and the older whole upload.
204#[allow(clippy::too_many_arguments)]
205async fn cache(
206 mut request: Request,
207 bucket: &Bucket,
208 services: &Services,
209 method: &str,
210 job: &str,
211 token: &str,
212 repo: &str,
213 what: &str,
214) -> Result<Response> {
215 let parts: Vec<&str> = what.split('/').collect();
216 let upload_id = query(&request, "upload").unwrap_or_default();
Merge branch 'worktree-agent-a3abfcce648e87dca'217 // The hash of the entry's paths and compression; runners from before
218 // it was sent send none.
219 let version = query(&request, "version").unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents220 match (method, parts.as_slice()) {
221 ("GET", ["cache"]) => {
222 let key = query(&request, "key").unwrap_or_default();
223 let restore: Vec<String> =
224 query(&request, "restore").unwrap_or_default().lines().map(str::trim).filter(|p| !p.is_empty()).map(str::to_owned).collect();
225 let found: Outcome<Option<CacheHit>> = g1t_kit::call(
226 &services.actions,
227 "cache_lookup",
Merge branch 'worktree-agent-a3abfcce648e87dca'228 &CacheLookupArgs { job: job.to_owned(), token: token.to_owned(), key: key.clone(), restore, version: Some(version.clone()).filter(|v| !v.is_empty()) },
Fast pages, required checks on the branch, self-hosted runners, honest incidents229 )
230 .await?;
231 let found = match refused(found) {
232 Ok(found) => found,
233 Err(reply) => return reply,
234 };
235 if let Some(hit) = found
236 && let Some(object) = bucket.get(&hit.object).execute().await?
237 && let Some(body) = object.body()
238 {
239 let mut response = Response::from_body(body.response_body()?)?;
240 let headers = response.headers_mut();
241 headers.set("x-g1t-key", &hit.key)?;
242 headers.set("content-length", &object.size().to_string())?;
243 headers.set("content-type", "application/octet-stream")?;
244 return Ok(response);
245 }
Merge branch 'worktree-agent-a3abfcce648e87dca'246 // Entries kept in KV before the cache moved to R2 had no scope,
247 // so they are never restored.
248 error(404, "Nothing cached under those keys.")
Fast pages, required checks on the branch, self-hosted runners, honest incidents249 }
250 // Older runners send a whole entry of at most 60 MB at once.
251 ("PUT", ["cache"]) => {
A repository has its own sidebar, as settings do252 let key = query(&request, "key").unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents253 let bytes = request.bytes().await?;
254 if bytes.len() > MAX_BYTES {
255 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 do256 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents257 let reserved: Outcome<CacheReservation> = g1t_kit::call(
258 &services.actions,
259 "cache_reserve",
Merge branch 'worktree-agent-a3abfcce648e87dca'260 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size: bytes.len() as u64, version: Some(version.clone()).filter(|v| !v.is_empty()) },
Fast pages, required checks on the branch, self-hosted runners, honest incidents261 )
262 .await?;
263 let reserved = match reserved {
264 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
265 return crate::reply(&json!({ "saved": false, "reason": failure.message }));
A repository has its own sidebar, as settings do266 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents267 other => match refused(other) {
268 Ok(reserved) => reserved,
269 Err(reply) => return reply,
270 },
271 };
272 let size = bytes.len() as u64;
273 bucket.put(&reserved.object, bytes).execute().await?;
274 commit(bucket, services, job, token, &reserved.id, size).await?;
275 crate::reply(&json!({ "saved": true }))
276 }
277 ("POST", ["cache", "uploads"]) => {
278 let key = query(&request, "key").unwrap_or_default();
279 let size = query(&request, "size").and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
280 let reserved: Outcome<CacheReservation> = g1t_kit::call(
281 &services.actions,
282 "cache_reserve",
Merge branch 'worktree-agent-a3abfcce648e87dca'283 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size, version: Some(version.clone()).filter(|v| !v.is_empty()) },
Fast pages, required checks on the branch, self-hosted runners, honest incidents284 )
285 .await?;
286 let reserved = match refused(reserved) {
287 Ok(reserved) => reserved,
288 Err(reply) => return reply,
289 };
290 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
291 crate::reply(&json!({ "id": reserved.id, "upload": upload.upload_id().await, "part_bytes": CACHE_PART_BYTES }))
292 }
293 ("PUT", ["cache", "uploads", id, part]) => {
294 let part = part.parse::<u16>().unwrap_or(0);
295 if part == 0 || upload_id.is_empty() {
296 return error(400, "A part is numbered from 1, and names its upload.");
A repository has its own sidebar, as settings do297 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents298 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
299 if length == 0 || length > CACHE_PART_BYTES {
300 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 do301 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents302 // The body goes to R2 as it comes, never held whole here.
303 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
304 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
305 let uploaded = upload.upload_part(part, body).await?;
306 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
307 }
308 ("POST", ["cache", "uploads", id, "complete"]) => {
309 let done: Complete = match request.json().await {
310 Ok(done) => done,
311 Err(_) => return error(400, "Send { size, parts: [{ part, etag }] }."),
312 };
313 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
314 let mut parts = done.parts;
315 parts.sort_by_key(|p| p.part);
316 if let Err(problem) = upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
317 let _ = abort(services, job, token, id).await;
318 return error(400, &format!("The upload could not be completed: {problem}"));
319 }
320 commit(bucket, services, job, token, id, done.size).await?;
321 crate::reply(&json!({ "saved": true }))
322 }
323 ("DELETE", ["cache", "uploads", id]) => {
324 if let Ok(upload) = bucket.resume_multipart_upload(object_of(repo, id), &upload_id) {
325 let _ = upload.abort().await;
A repository has its own sidebar, as settings do326 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents327 abort(services, job, token, id).await?;
328 crate::reply(&json!({ "aborted": true }))
A repository has its own sidebar, as settings do329 }
330 _ => error(404, "No such endpoint."),
331 }
332}
333
Fast pages, required checks on the branch, self-hosted runners, honest incidents334/// Where an entry is in R2: under its repository, by its id, as the
335/// actions service named it when it was reserved.
336fn object_of(repo: &str, id: &str) -> String {
337 format!("c/{repo}/{id}")
338}
339
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2340/// Where an artifact is in R2, as the actions service names it.
341fn artifact_object(repo: &str, id: u64) -> String {
342 format!("a/{repo}/{id}")
343}
344
345/// An artifact as a job's runner lists it.
346fn for_runner(artifact: &Artifact) -> Value {
347 json!({
348 "id": artifact.id,
349 "name": artifact.name,
350 "size": artifact.size,
351 "digest": artifact.digest,
352 "format": artifact.format,
353 "created_at": artifact.created_at,
354 "expires_at": artifact.expires_at,
355 })
356}
357
358/// An R2 object streamed back, with what it is.
359async fn stream(bucket: &Bucket, object: &str, format: &str) -> Result<Response> {
360 let Some(found) = bucket.get(object).execute().await? else { return error(404, "That artifact is gone: it expired or was deleted.") };
361 let size = found.size();
362 let Some(body) = found.body() else { return error(404, "That artifact is gone: it expired or was deleted.") };
363 let mut response = Response::from_body(body.response_body()?)?;
364 let headers = response.headers_mut();
365 headers.set("content-length", &size.to_string())?;
366 headers.set("content-type", if format == "tgz" { "application/gzip" } else { "application/zip" })?;
367 headers.set("x-g1t-format", format)?;
368 Ok(response)
369}
370
371/// A job's artifacts: listing its run's (or another run's of its
372/// repository), uploading in parts, downloading, deleting; and the whole
373/// uploads and downloads by name of older runners.
374#[allow(clippy::too_many_arguments)]
375async fn artifacts(
376 mut request: Request,
377 kv: &KvStore,
378 bucket: &Bucket,
379 services: &Services,
380 method: &str,
381 job: &str,
382 token: &str,
383 run: &str,
384 repo: &str,
385 what: &str,
386) -> Result<Response> {
387 let parts: Vec<&str> = what.split('/').collect();
388 let upload_id = query(&request, "upload").unwrap_or_default();
389 let credential = |id: Option<u64>, name: Option<String>, run_id: Option<String>| JobArtifactsArgs {
390 job: job.to_owned(),
391 token: token.to_owned(),
392 run_id,
393 name,
394 id,
395 };
396 let actions = &services.actions;
397 match (method, parts.as_slice()) {
398 ("GET", ["artifacts"]) => {
399 let run_id = query(&request, "run_id").filter(|r| !r.is_empty());
400 let found: Outcome<Vec<Artifact>> = g1t_kit::call(actions, "job_artifacts", &credential(None, None, run_id.clone())).await?;
401 match refused(found) {
402 Ok(found) => {
403 let mut listed: Vec<Value> = found.iter().map(for_runner).collect();
404 // Artifacts older runners kept in KV, for the days they
405 // are still there.
406 let legacy_run = run_id.as_deref().unwrap_or(run);
407 for (_, meta) in list(kv, &format!("a/{legacy_run}/")).await? {
408 if !listed.iter().any(|a| a["name"] == meta.name.as_str()) {
409 listed.push(json!({ "name": meta.name, "size": meta.size, "format": "tgz" }));
410 }
411 }
412 crate::reply(&listed)
413 }
414 Err(reply) => reply,
415 }
416 }
417 ("POST", ["artifacts", "uploads"]) => {
418 let flag = |name: &str| query(&request, name).is_some_and(|v| v == "true");
419 let args = ArtifactReserveArgs {
420 job: job.to_owned(),
421 token: token.to_owned(),
422 name: query(&request, "name").unwrap_or_default(),
423 size: query(&request, "size").and_then(|s| s.parse().ok()).unwrap_or(0),
424 retention_days: query(&request, "retention_days").and_then(|d| d.parse().ok()).unwrap_or(0),
425 expires_at: None,
426 overwrite: flag("overwrite"),
427 format: query(&request, "format"),
428 };
429 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
430 let reserved = match refused(reserved) {
431 Ok(reserved) => reserved,
432 Err(reply) => return reply,
433 };
434 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
435 crate::reply(&json!({
436 "id": reserved.id,
437 "upload": upload.upload_id().await,
438 "part_bytes": ARTIFACT_PART_BYTES,
439 "retention_days": reserved.retention_days,
440 "expires_at": reserved.expires_at,
441 }))
442 }
443 ("PUT", ["artifacts", "uploads", id, part]) => {
444 let (Ok(id), Ok(part)) = (id.parse::<u64>(), part.parse::<u16>()) else {
445 return error(400, "A part is numbered from 1, of an artifact named by its number.");
446 };
447 if part == 0 || upload_id.is_empty() {
448 return error(400, "A part is numbered from 1, and names its upload.");
449 }
450 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
451 if length == 0 || length > ARTIFACT_PART_BYTES {
452 return error(413, &format!("A part is 1 to {} MB, with its length.", ARTIFACT_PART_BYTES / 1_048_576));
453 }
454 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
455 let upload = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id)?;
456 let uploaded = upload.upload_part(part, body).await?;
457 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
458 }
459 ("POST", ["artifacts", "uploads", id, "complete"]) => {
460 let Ok(id) = id.parse::<u64>() else { return error(404, "No such upload.") };
461 let done: Value = request.json().await.unwrap_or(Value::Null);
462 let mut parts: Vec<Part> = serde_json::from_value(done["parts"].clone()).unwrap_or_default();
463 if parts.is_empty() {
464 return error(400, "Send { size, parts: [{ part, etag }], digest }.");
465 }
466 parts.sort_by_key(|p| p.part);
467 let upload = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id)?;
468 let object = match upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
469 Ok(object) => object,
470 Err(problem) => {
471 let _: Outcome<bool> = g1t_kit::call(actions, "artifact_abort", &credential(Some(id), None, None)).await?;
472 return error(400, &format!("The upload could not be completed: {problem}"));
473 }
474 };
475 let args = ArtifactCommitArgs {
476 job: job.to_owned(),
477 token: token.to_owned(),
478 id: Some(id),
479 name: None,
480 // What R2 holds, not what the runner says.
481 size: object.size(),
482 digest: done["digest"].as_str().map(str::to_owned),
483 };
484 let committed: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
485 match refused(committed) {
486 Ok(artifact) => crate::reply(&for_runner(&artifact)),
487 Err(reply) => reply,
488 }
489 }
490 ("DELETE", ["artifacts", "uploads", id]) => {
491 let Ok(id) = id.parse::<u64>() else { return error(404, "No such upload.") };
492 if let Ok(upload) = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id) {
493 let _ = upload.abort().await;
494 }
495 let _: Outcome<bool> = g1t_kit::call(actions, "artifact_abort", &credential(Some(id), None, None)).await?;
496 crate::reply(&json!({ "aborted": true }))
497 }
498 ("GET", ["artifacts", id, "download"]) => {
499 let Ok(id) = id.parse::<u64>() else { return error(404, "No such artifact.") };
500 // Any run of the job's repository: the service checks.
501 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &credential(Some(id), None, query(&request, "run_id"))).await?;
502 match found {
503 Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await,
504 Outcome::Fail(failure) => error(failure.code.http_status(), &failure.message),
505 }
506 }
507 ("DELETE", ["artifacts", id]) => {
508 let Ok(id) = id.parse::<u64>() else { return error(404, "No such artifact.") };
509 let done: Outcome<Artifact> = g1t_kit::call(actions, "job_delete_artifact", &credential(Some(id), None, None)).await?;
510 match refused(done) {
511 Ok(artifact) => crate::reply(&for_runner(&artifact)),
512 Err(reply) => reply,
513 }
514 }
515 // Older runners: a whole artifact of at most 60 MB, by name.
516 ("PUT", ["artifacts", name]) => {
517 let name = decode(name);
518 let bytes = request.bytes().await?;
519 if bytes.len() > MAX_BYTES {
520 return error(413, "An artifact sent at once is at most 60 MB; newer runners upload it in parts.");
521 }
522 let args = ArtifactReserveArgs {
523 job: job.to_owned(),
524 token: token.to_owned(),
525 name: name.clone(),
526 size: bytes.len() as u64,
527 format: Some("tgz".to_owned()),
528 ..ArtifactReserveArgs::default()
529 };
530 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
531 let reserved = match refused(reserved) {
532 Ok(reserved) => reserved,
533 Err(reply) => return reply,
534 };
535 let size = bytes.len() as u64;
536 bucket.put(&reserved.object, bytes).execute().await?;
537 let args = ArtifactCommitArgs { job: job.to_owned(), token: token.to_owned(), id: Some(reserved.id), name: None, size, digest: None };
538 let committed: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
539 match refused(committed) {
540 Ok(_) => crate::reply(&json!({ "name": name, "size": size })),
541 Err(reply) => reply,
542 }
543 }
544 ("GET", ["artifacts", name]) => {
545 let name = decode(name);
546 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &credential(None, Some(name.clone()), None)).await?;
547 match found {
548 Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await,
549 // One an older runner kept in KV.
550 Outcome::Fail(_) => match get(kv, &format!("a/{run}/{name}")).await? {
551 Some(bytes) => Response::from_bytes(bytes),
552 None => error(404, "No such artifact."),
553 },
554 }
555 }
556 _ => error(404, "No such endpoint."),
557 }
558}
559
Fast pages, required checks on the branch, self-hosted runners, honest incidents560/// Marks an uploaded entry ready, and deletes what that evicted.
561async fn commit(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<()> {
562 let committed: Outcome<CacheCommitted> = g1t_kit::call(
563 &services.actions,
564 "cache_commit",
565 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
566 )
567 .await?;
568 if let Outcome::Ok(committed) = committed
569 && !committed.evicted.is_empty()
570 {
571 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
572 }
573 Ok(())
574}
575
576async fn abort(services: &Services, job: &str, token: &str, id: &str) -> Result<()> {
577 let _: Outcome<bool> = g1t_kit::call(
578 &services.actions,
579 "cache_abort",
580 &CacheAbortArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned() },
581 )
582 .await?;
583 Ok(())
584}
585
586/// An entry saved in KV before the cache moved to R2, by key or restore key.
A repository has its own sidebar, as settings do587/// Someone who can see the run downloading one of its artifacts.
588pub async fn download(env: &Env, services: &Services, viewer: &g1t_contracts::Viewer, owner: &str, repo: &str, run: &str, name: &str) -> Result<Response> {
589 let seen: Outcome<Value> = g1t_kit::call(
590 &services.actions,
591 "run",
592 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
593 )
594 .await?;
595 if let Outcome::Fail(refused) = seen {
596 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
597 return error(status, &refused.message);
598 }
599 let name = decode(name);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2600 // Kept in R2: a redirect to a link signed for a few minutes.
601 let args = ArtifactArgs {
602 repo: g1t_contracts::repos::RepoPath { namespace: owner.to_owned(), name: repo.to_owned() },
603 viewer: viewer.clone(),
604 id: None,
605 run: Some(run.to_owned()),
606 name: Some(name.clone()),
607 };
608 let found: Outcome<ArtifactBlob> = g1t_kit::call(&services.actions, "artifact_download", &args).await?;
609 if let Outcome::Ok(found) = found
610 && !found.blob.is_empty()
611 {
612 return Response::redirect_with_status(Url::parse(&crate::artifacts::blob_url(&services.addresses.api, &found.blob))?, 302);
613 }
614 // Kept in KV by an older runner.
A repository has its own sidebar, as settings do615 match get(&store(env)?, &format!("a/{run}/{name}")).await? {
616 Some(bytes) => {
617 let mut response = Response::from_bytes(bytes)?;
618 let headers = response.headers_mut();
619 headers.set("content-type", "application/gzip")?;
620 headers.set("content-disposition", &format!("attachment; filename=\"{}.tar.gz\"", name.replace('"', "")))?;
621 Ok(response)
622 }
623 None => error(404, "No such artifact, or it has expired."),
624 }
625}
626

This file's history is long; its oldest lines are credited to the oldest commit read.