Skip to content
642 linesCodeBlameRaw
1//! 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.)
4//!
5//! 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).
12//!
13//! A sandbox reaches these with its job's token:
14//!
15//! - `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)
22//! - `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)
28//!
29//! People download an artifact through the REST API (artifacts.rs), or by
30//! name at `/repos/{owner}/{repo}/actions/runs/{run}/artifacts/{name}`.
31
32use serde::{Deserialize, Serialize};
33use serde_json::{Value, json};
34use worker::kv::KvStore;
35use worker::{Bucket, Env, Request, Response, Result, UploadedPart, Url};
36
37use g1t_contracts::actions::{
38 ARTIFACT_PART_BYTES, Artifact, ArtifactArgs, ArtifactBlob, ArtifactCommitArgs, ArtifactReservation, ArtifactReserveArgs, CACHE_PART_BYTES,
39 CacheAbortArgs, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation, CacheReserveArgs, JobArtifactsArgs,
40};
41use g1t_contracts::{FailureCode, Outcome};
42
43use crate::operations::Services;
44
45/// The largest artifact or cache entry an older runner sends at once,
46/// held in a Worker's memory.
47const 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> {
104 Ok(crate::reply(&json!({ "error": { "message": message } }))?.with_status(status))
105}
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/`.
145pub async fn for_job(request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
146 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) {
161 (_, 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
164 }
165 (_, what) if what == "cache" || what.starts_with("cache/") => {
166 let bucket = env.bucket("ACTIONS_CACHE")?;
167 cache(request, &kv, &bucket, services, method, job, &token, &repo, what).await
168 }
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 kv: &KvStore,
208 bucket: &Bucket,
209 services: &Services,
210 method: &str,
211 job: &str,
212 token: &str,
213 repo: &str,
214 what: &str,
215) -> Result<Response> {
216 let parts: Vec<&str> = what.split('/').collect();
217 let upload_id = query(&request, "upload").unwrap_or_default();
218 match (method, parts.as_slice()) {
219 ("GET", ["cache"]) => {
220 let key = query(&request, "key").unwrap_or_default();
221 let restore: Vec<String> =
222 query(&request, "restore").unwrap_or_default().lines().map(str::trim).filter(|p| !p.is_empty()).map(str::to_owned).collect();
223 let found: Outcome<Option<CacheHit>> = g1t_kit::call(
224 &services.actions,
225 "cache_lookup",
226 &CacheLookupArgs { job: job.to_owned(), token: token.to_owned(), key: key.clone(), restore: restore.clone(), version: None },
227 )
228 .await?;
229 let found = match refused(found) {
230 Ok(found) => found,
231 Err(reply) => return reply,
232 };
233 if let Some(hit) = found
234 && let Some(object) = bucket.get(&hit.object).execute().await?
235 && let Some(body) = object.body()
236 {
237 let mut response = Response::from_body(body.response_body()?)?;
238 let headers = response.headers_mut();
239 headers.set("x-g1t-key", &hit.key)?;
240 headers.set("content-length", &object.size().to_string())?;
241 headers.set("content-type", "application/octet-stream")?;
242 return Ok(response);
243 }
244 // Entries saved in KV before the cache moved to R2.
245 kv_lookup(kv, repo, &key, &restore).await
246 }
247 // Older runners send a whole entry of at most 60 MB at once.
248 ("PUT", ["cache"]) => {
249 let key = query(&request, "key").unwrap_or_default();
250 let bytes = request.bytes().await?;
251 if bytes.len() > MAX_BYTES {
252 return error(413, "An entry sent at once is at most 60 MB; newer runners upload it in parts.");
253 }
254 let reserved: Outcome<CacheReservation> = g1t_kit::call(
255 &services.actions,
256 "cache_reserve",
257 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size: bytes.len() as u64, version: None },
258 )
259 .await?;
260 let reserved = match reserved {
261 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
262 return crate::reply(&json!({ "saved": false, "reason": failure.message }));
263 }
264 other => match refused(other) {
265 Ok(reserved) => reserved,
266 Err(reply) => return reply,
267 },
268 };
269 let size = bytes.len() as u64;
270 bucket.put(&reserved.object, bytes).execute().await?;
271 commit(bucket, services, job, token, &reserved.id, size).await?;
272 crate::reply(&json!({ "saved": true }))
273 }
274 ("POST", ["cache", "uploads"]) => {
275 let key = query(&request, "key").unwrap_or_default();
276 let size = query(&request, "size").and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
277 let reserved: Outcome<CacheReservation> = g1t_kit::call(
278 &services.actions,
279 "cache_reserve",
280 &CacheReserveArgs { job: job.to_owned(), token: token.to_owned(), key, size, version: None },
281 )
282 .await?;
283 let reserved = match refused(reserved) {
284 Ok(reserved) => reserved,
285 Err(reply) => return reply,
286 };
287 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
288 crate::reply(&json!({ "id": reserved.id, "upload": upload.upload_id().await, "part_bytes": CACHE_PART_BYTES }))
289 }
290 ("PUT", ["cache", "uploads", id, part]) => {
291 let part = part.parse::<u16>().unwrap_or(0);
292 if part == 0 || upload_id.is_empty() {
293 return error(400, "A part is numbered from 1, and names its upload.");
294 }
295 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
296 if length == 0 || length > CACHE_PART_BYTES {
297 return error(413, &format!("A part is 1 to {} MB, with its length.", CACHE_PART_BYTES / 1_048_576));
298 }
299 // The body goes to R2 as it comes, never held whole here.
300 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
301 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
302 let uploaded = upload.upload_part(part, body).await?;
303 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
304 }
305 ("POST", ["cache", "uploads", id, "complete"]) => {
306 let done: Complete = match request.json().await {
307 Ok(done) => done,
308 Err(_) => return error(400, "Send { size, parts: [{ part, etag }] }."),
309 };
310 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
311 let mut parts = done.parts;
312 parts.sort_by_key(|p| p.part);
313 if let Err(problem) = upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
314 let _ = abort(services, job, token, id).await;
315 return error(400, &format!("The upload could not be completed: {problem}"));
316 }
317 commit(bucket, services, job, token, id, done.size).await?;
318 crate::reply(&json!({ "saved": true }))
319 }
320 ("DELETE", ["cache", "uploads", id]) => {
321 if let Ok(upload) = bucket.resume_multipart_upload(object_of(repo, id), &upload_id) {
322 let _ = upload.abort().await;
323 }
324 abort(services, job, token, id).await?;
325 crate::reply(&json!({ "aborted": true }))
326 }
327 _ => error(404, "No such endpoint."),
328 }
329}
330
331/// Where an entry is in R2: under its repository, by its id, as the
332/// actions service named it when it was reserved.
333fn object_of(repo: &str, id: &str) -> String {
334 format!("c/{repo}/{id}")
335}
336
337/// Where an artifact is in R2, as the actions service names it.
338fn artifact_object(repo: &str, id: u64) -> String {
339 format!("a/{repo}/{id}")
340}
341
342/// An artifact as a job's runner lists it.
343fn for_runner(artifact: &Artifact) -> Value {
344 json!({
345 "id": artifact.id,
346 "name": artifact.name,
347 "size": artifact.size,
348 "digest": artifact.digest,
349 "format": artifact.format,
350 "created_at": artifact.created_at,
351 "expires_at": artifact.expires_at,
352 })
353}
354
355/// An R2 object streamed back, with what it is.
356async fn stream(bucket: &Bucket, object: &str, format: &str) -> Result<Response> {
357 let Some(found) = bucket.get(object).execute().await? else { return error(404, "That artifact is gone: it expired or was deleted.") };
358 let size = found.size();
359 let Some(body) = found.body() else { return error(404, "That artifact is gone: it expired or was deleted.") };
360 let mut response = Response::from_body(body.response_body()?)?;
361 let headers = response.headers_mut();
362 headers.set("content-length", &size.to_string())?;
363 headers.set("content-type", if format == "tgz" { "application/gzip" } else { "application/zip" })?;
364 headers.set("x-g1t-format", format)?;
365 Ok(response)
366}
367
368/// A job's artifacts: listing its run's (or another run's of its
369/// repository), uploading in parts, downloading, deleting; and the whole
370/// uploads and downloads by name of older runners.
371#[allow(clippy::too_many_arguments)]
372async fn artifacts(
373 mut request: Request,
374 kv: &KvStore,
375 bucket: &Bucket,
376 services: &Services,
377 method: &str,
378 job: &str,
379 token: &str,
380 run: &str,
381 repo: &str,
382 what: &str,
383) -> Result<Response> {
384 let parts: Vec<&str> = what.split('/').collect();
385 let upload_id = query(&request, "upload").unwrap_or_default();
386 let credential = |id: Option<u64>, name: Option<String>, run_id: Option<String>| JobArtifactsArgs {
387 job: job.to_owned(),
388 token: token.to_owned(),
389 run_id,
390 name,
391 id,
392 };
393 let actions = &services.actions;
394 match (method, parts.as_slice()) {
395 ("GET", ["artifacts"]) => {
396 let run_id = query(&request, "run_id").filter(|r| !r.is_empty());
397 let found: Outcome<Vec<Artifact>> = g1t_kit::call(actions, "job_artifacts", &credential(None, None, run_id.clone())).await?;
398 match refused(found) {
399 Ok(found) => {
400 let mut listed: Vec<Value> = found.iter().map(for_runner).collect();
401 // Artifacts older runners kept in KV, for the days they
402 // are still there.
403 let legacy_run = run_id.as_deref().unwrap_or(run);
404 for (_, meta) in list(kv, &format!("a/{legacy_run}/")).await? {
405 if !listed.iter().any(|a| a["name"] == meta.name.as_str()) {
406 listed.push(json!({ "name": meta.name, "size": meta.size, "format": "tgz" }));
407 }
408 }
409 crate::reply(&listed)
410 }
411 Err(reply) => reply,
412 }
413 }
414 ("POST", ["artifacts", "uploads"]) => {
415 let flag = |name: &str| query(&request, name).is_some_and(|v| v == "true");
416 let args = ArtifactReserveArgs {
417 job: job.to_owned(),
418 token: token.to_owned(),
419 name: query(&request, "name").unwrap_or_default(),
420 size: query(&request, "size").and_then(|s| s.parse().ok()).unwrap_or(0),
421 retention_days: query(&request, "retention_days").and_then(|d| d.parse().ok()).unwrap_or(0),
422 expires_at: None,
423 overwrite: flag("overwrite"),
424 format: query(&request, "format"),
425 };
426 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
427 let reserved = match refused(reserved) {
428 Ok(reserved) => reserved,
429 Err(reply) => return reply,
430 };
431 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
432 crate::reply(&json!({
433 "id": reserved.id,
434 "upload": upload.upload_id().await,
435 "part_bytes": ARTIFACT_PART_BYTES,
436 "retention_days": reserved.retention_days,
437 "expires_at": reserved.expires_at,
438 }))
439 }
440 ("PUT", ["artifacts", "uploads", id, part]) => {
441 let (Ok(id), Ok(part)) = (id.parse::<u64>(), part.parse::<u16>()) else {
442 return error(400, "A part is numbered from 1, of an artifact named by its number.");
443 };
444 if part == 0 || upload_id.is_empty() {
445 return error(400, "A part is numbered from 1, and names its upload.");
446 }
447 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
448 if length == 0 || length > ARTIFACT_PART_BYTES {
449 return error(413, &format!("A part is 1 to {} MB, with its length.", ARTIFACT_PART_BYTES / 1_048_576));
450 }
451 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
452 let upload = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id)?;
453 let uploaded = upload.upload_part(part, body).await?;
454 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
455 }
456 ("POST", ["artifacts", "uploads", id, "complete"]) => {
457 let Ok(id) = id.parse::<u64>() else { return error(404, "No such upload.") };
458 let done: Value = request.json().await.unwrap_or(Value::Null);
459 let mut parts: Vec<Part> = serde_json::from_value(done["parts"].clone()).unwrap_or_default();
460 if parts.is_empty() {
461 return error(400, "Send { size, parts: [{ part, etag }], digest }.");
462 }
463 parts.sort_by_key(|p| p.part);
464 let upload = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id)?;
465 let object = match upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
466 Ok(object) => object,
467 Err(problem) => {
468 let _: Outcome<bool> = g1t_kit::call(actions, "artifact_abort", &credential(Some(id), None, None)).await?;
469 return error(400, &format!("The upload could not be completed: {problem}"));
470 }
471 };
472 let args = ArtifactCommitArgs {
473 job: job.to_owned(),
474 token: token.to_owned(),
475 id: Some(id),
476 name: None,
477 // What R2 holds, not what the runner says.
478 size: object.size(),
479 digest: done["digest"].as_str().map(str::to_owned),
480 };
481 let committed: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
482 match refused(committed) {
483 Ok(artifact) => crate::reply(&for_runner(&artifact)),
484 Err(reply) => reply,
485 }
486 }
487 ("DELETE", ["artifacts", "uploads", id]) => {
488 let Ok(id) = id.parse::<u64>() else { return error(404, "No such upload.") };
489 if let Ok(upload) = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id) {
490 let _ = upload.abort().await;
491 }
492 let _: Outcome<bool> = g1t_kit::call(actions, "artifact_abort", &credential(Some(id), None, None)).await?;
493 crate::reply(&json!({ "aborted": true }))
494 }
495 ("GET", ["artifacts", id, "download"]) => {
496 let Ok(id) = id.parse::<u64>() else { return error(404, "No such artifact.") };
497 // Any run of the job's repository: the service checks.
498 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &credential(Some(id), None, query(&request, "run_id"))).await?;
499 match found {
500 Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await,
501 Outcome::Fail(failure) => error(failure.code.http_status(), &failure.message),
502 }
503 }
504 ("DELETE", ["artifacts", id]) => {
505 let Ok(id) = id.parse::<u64>() else { return error(404, "No such artifact.") };
506 let done: Outcome<Artifact> = g1t_kit::call(actions, "job_delete_artifact", &credential(Some(id), None, None)).await?;
507 match refused(done) {
508 Ok(artifact) => crate::reply(&for_runner(&artifact)),
509 Err(reply) => reply,
510 }
511 }
512 // Older runners: a whole artifact of at most 60 MB, by name.
513 ("PUT", ["artifacts", name]) => {
514 let name = decode(name);
515 let bytes = request.bytes().await?;
516 if bytes.len() > MAX_BYTES {
517 return error(413, "An artifact sent at once is at most 60 MB; newer runners upload it in parts.");
518 }
519 let args = ArtifactReserveArgs {
520 job: job.to_owned(),
521 token: token.to_owned(),
522 name: name.clone(),
523 size: bytes.len() as u64,
524 format: Some("tgz".to_owned()),
525 ..ArtifactReserveArgs::default()
526 };
527 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
528 let reserved = match refused(reserved) {
529 Ok(reserved) => reserved,
530 Err(reply) => return reply,
531 };
532 let size = bytes.len() as u64;
533 bucket.put(&reserved.object, bytes).execute().await?;
534 let args = ArtifactCommitArgs { job: job.to_owned(), token: token.to_owned(), id: Some(reserved.id), name: None, size, digest: None };
535 let committed: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
536 match refused(committed) {
537 Ok(_) => crate::reply(&json!({ "name": name, "size": size })),
538 Err(reply) => reply,
539 }
540 }
541 ("GET", ["artifacts", name]) => {
542 let name = decode(name);
543 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &credential(None, Some(name.clone()), None)).await?;
544 match found {
545 Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await,
546 // One an older runner kept in KV.
547 Outcome::Fail(_) => match get(kv, &format!("a/{run}/{name}")).await? {
548 Some(bytes) => Response::from_bytes(bytes),
549 None => error(404, "No such artifact."),
550 },
551 }
552 }
553 _ => error(404, "No such endpoint."),
554 }
555}
556
557/// Marks an uploaded entry ready, and deletes what that evicted.
558async fn commit(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<()> {
559 let committed: Outcome<CacheCommitted> = g1t_kit::call(
560 &services.actions,
561 "cache_commit",
562 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
563 )
564 .await?;
565 if let Outcome::Ok(committed) = committed
566 && !committed.evicted.is_empty()
567 {
568 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
569 }
570 Ok(())
571}
572
573async fn abort(services: &Services, job: &str, token: &str, id: &str) -> Result<()> {
574 let _: Outcome<bool> = g1t_kit::call(
575 &services.actions,
576 "cache_abort",
577 &CacheAbortArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned() },
578 )
579 .await?;
580 Ok(())
581}
582
583/// An entry saved in KV before the cache moved to R2, by key or restore key.
584async fn kv_lookup(kv: &KvStore, repo: &str, key: &str, restore: &[String]) -> Result<Response> {
585 // The exact key, else the newest entry under each restore key.
586 if let Some(bytes) = get(kv, &format!("c/{repo}/{key}")).await? {
587 let mut response = Response::from_bytes(bytes)?;
588 response.headers_mut().set("x-g1t-key", key)?;
589 return Ok(response);
590 }
591 for prefix in restore {
592 if let Some((base, meta)) = list(kv, &format!("c/{repo}/{prefix}")).await?.into_iter().next()
593 && let Some(bytes) = get(kv, &base).await?
594 {
595 let mut response = Response::from_bytes(bytes)?;
596 response.headers_mut().set("x-g1t-key", &meta.name)?;
597 return Ok(response);
598 }
599 }
600 error(404, "Nothing cached under those keys.")
601}
602
603/// Someone who can see the run downloading one of its artifacts.
604pub async fn download(env: &Env, services: &Services, viewer: &g1t_contracts::Viewer, owner: &str, repo: &str, run: &str, name: &str) -> Result<Response> {
605 let seen: Outcome<Value> = g1t_kit::call(
606 &services.actions,
607 "run",
608 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
609 )
610 .await?;
611 if let Outcome::Fail(refused) = seen {
612 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
613 return error(status, &refused.message);
614 }
615 let name = decode(name);
616 // Kept in R2: a redirect to a link signed for a few minutes.
617 let args = ArtifactArgs {
618 repo: g1t_contracts::repos::RepoPath { namespace: owner.to_owned(), name: repo.to_owned() },
619 viewer: viewer.clone(),
620 id: None,
621 run: Some(run.to_owned()),
622 name: Some(name.clone()),
623 };
624 let found: Outcome<ArtifactBlob> = g1t_kit::call(&services.actions, "artifact_download", &args).await?;
625 if let Outcome::Ok(found) = found
626 && !found.blob.is_empty()
627 {
628 return Response::redirect_with_status(Url::parse(&crate::artifacts::blob_url(&services.addresses.api, &found.blob))?, 302);
629 }
630 // Kept in KV by an older runner.
631 match get(&store(env)?, &format!("a/{run}/{name}")).await? {
632 Some(bytes) => {
633 let mut response = Response::from_bytes(bytes)?;
634 let headers = response.headers_mut();
635 headers.set("content-type", "application/gzip")?;
636 headers.set("content-disposition", &format!("attachment; filename=\"{}.tar.gz\"", name.replace('"', "")))?;
637 Ok(response)
638 }
639 None => error(404, "No such artifact, or it has expired."),
640 }
641}
642