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 R2 | 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.) | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 4 | //! |
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 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). | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 12 | //! |
| 13 | //! A sandbox reaches these with its job's token: | |
| 14 | //! | |
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 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) | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 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) | |
| A repository has its own sidebar, as settings do | 28 | //! |
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 29 | //! 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 do | 31 | |
| 32 | use serde::{Deserialize, Serialize}; | |
| 33 | use serde_json::{Value, json}; | |
| 34 | use worker::kv::KvStore; | |
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 35 | use worker::{Bucket, Env, Request, Response, Result, UploadedPart, Url}; |
| A repository has its own sidebar, as settings do | 36 | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 37 | use g1t_contracts::actions::{ |
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 38 | 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 incidents | 40 | }; |
| A repository has its own sidebar, as settings do | 41 | use g1t_contracts::{FailureCode, Outcome}; |
| 42 | ||
| 43 | use crate::operations::Services; | |
| 44 | ||
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 45 | /// 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 do | 47 | const MAX_BYTES: usize = 60 * 1024 * 1024; |
| 48 | ||
| 49 | #[derive(Serialize, Deserialize)] | |
| 50 | struct 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 | ||
| 59 | fn store(env: &Env) -> Result<KvStore> { | |
| 60 | env.kv("BLOBS") | |
| 61 | } | |
| 62 | ||
| 63 | async 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. | |
| 77 | async 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 | ||
| 103 | fn 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 API | 104 | Ok(crate::reply(&json!({ "error": { "message": message } }))?.with_status(status)) |
| A repository has its own sidebar, as settings do | 105 | } |
| 106 | ||
| 107 | fn 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 | ||
| 138 | fn 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 R2 | 145 | pub async fn for_job(request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> { |
| A repository has its own sidebar, as settings do | 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) { | |
| Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2 | 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 | |
| A repository has its own sidebar, as settings do | 164 | } |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 165 | (_, 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 incidents | 168 | } |
| 169 | _ => error(404, "No such endpoint."), | |
| 170 | } | |
| 171 | } | |
| 172 | ||
| 173 | /// A part's number and the etag R2 gave it. | |
| 174 | #[derive(Deserialize)] | |
| 175 | struct Part { | |
| 176 | part: u16, | |
| 177 | etag: String, | |
| 178 | } | |
| 179 | ||
| 180 | #[derive(Deserialize)] | |
| 181 | struct 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. | |
| 187 | fn 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)] | |
| 205 | async 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 incidents | 220 | 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 incidents | 229 | ) |
| 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 incidents | 249 | } |
| 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 do | 252 | let key = query(&request, "key").unwrap_or_default(); |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 253 | 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 do | 256 | } |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 257 | 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 incidents | 261 | ) |
| 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 do | 266 | } |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 267 | 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 incidents | 284 | ) |
| 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 do | 297 | } |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 298 | 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 do | 301 | } |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 302 | // 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 do | 326 | } |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 327 | abort(services, job, token, id).await?; |
| 328 | crate::reply(&json!({ "aborted": true })) | |
| A repository has its own sidebar, as settings do | 329 | } |
| 330 | _ => error(404, "No such endpoint."), | |
| 331 | } | |
| 332 | } | |
| 333 | ||
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 334 | /// Where an entry is in R2: under its repository, by its id, as the |
| 335 | /// actions service named it when it was reserved. | |
| 336 | fn 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 R2 | 340 | /// Where an artifact is in R2, as the actions service names it. |
| 341 | fn artifact_object(repo: &str, id: u64) -> String { | |
| 342 | format!("a/{repo}/{id}") | |
| 343 | } | |
| 344 | ||
| 345 | /// An artifact as a job's runner lists it. | |
| 346 | fn 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. | |
| 359 | async 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)] | |
| 375 | async 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 incidents | 560 | /// Marks an uploaded entry ready, and deletes what that evicted. |
| 561 | async 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 | ||
| 576 | async 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 do | 587 | /// Someone who can see the run downloading one of its artifacts. |
| 588 | pub 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 R2 | 600 | // 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 do | 615 | 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.