Skip to content
655 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
Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails59/// Artifacts older runners kept in KV expire 14 days after they were made,
60/// and none has been made there since artifacts moved to R2 on 2026-10-08:
61/// from 2026-10-22T00:00Z every one is gone, and KV is not asked (a list is
62/// the dearest thing KV does). Delete the KV artifact code after that date,
63/// with its twin in apps/web/app/lib/artifacts.server.ts.
64const LEGACY_KV_UNTIL_MS: u64 = 1_792_627_200_000;
65
66/// Whether artifacts kept in KV may still be there at `now`.
67fn legacy_kv(now: u64) -> bool {
68 now < LEGACY_KV_UNTIL_MS
69}
70
A repository has its own sidebar, as settings do71fn store(env: &Env) -> Result<KvStore> {
72 env.kv("BLOBS")
73}
74
75async fn get(kv: &KvStore, base: &str) -> Result<Option<Vec<u8>>> {
76 let Some(meta) = kv.get(base).json::<Meta>().await? else { return Ok(None) };
77 let mut out = Vec::with_capacity(meta.size);
78 for index in 0..meta.chunks {
79 match kv.get(&format!("{base}#{index}")).bytes().await? {
80 Some(bytes) => out.extend_from_slice(&bytes),
81 // A chunk that expired first: the whole entry is gone.
82 None => return Ok(None),
83 }
84 }
85 Ok(Some(out))
86}
87
88/// The metadata of every entry under a prefix, newest first.
89async fn list(kv: &KvStore, prefix: &str) -> Result<Vec<(String, Meta)>> {
90 let mut found = Vec::new();
91 let mut cursor: Option<String> = None;
92 loop {
93 let mut listing = kv.list().prefix(prefix.to_owned());
94 if let Some(cursor) = cursor.take() {
95 listing = listing.cursor(cursor);
96 }
97 let page = listing.execute().await?;
98 for key in page.keys {
99 if key.name.contains('#') {
100 continue;
101 }
102 if let Some(meta) = key.metadata.and_then(|m| serde_json::from_value::<Meta>(m).ok()) {
103 found.push((key.name, meta));
104 }
105 }
106 if page.list_complete || page.cursor.is_none() {
107 break;
108 }
109 cursor = page.cursor;
110 }
111 found.sort_by_key(|entry| std::cmp::Reverse(entry.1.at));
112 Ok(found)
113}
114
115fn 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 API116 Ok(crate::reply(&json!({ "error": { "message": message } }))?.with_status(status))
A repository has its own sidebar, as settings do117}
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/`.
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2157pub async fn for_job(request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
A repository has its own sidebar, as settings do158 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) {
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2173 (_, what) if what == "artifacts" || what.starts_with("artifacts/") => {
174 let bucket = env.bucket("ACTIONS_CACHE")?;
175 artifacts(request, &kv, &bucket, services, method, job, &token, &run, &repo, what).await
A repository has its own sidebar, as settings do176 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents177 (_, what) if what == "cache" || what.starts_with("cache/") => {
178 let bucket = env.bucket("ACTIONS_CACHE")?;
Merge branch 'worktree-agent-a3abfcce648e87dca'179 cache(request, &bucket, services, method, job, &token, &repo, what).await
Fast pages, required checks on the branch, self-hosted runners, honest incidents180 }
181 _ => error(404, "No such endpoint."),
182 }
183}
184
185/// A part's number and the etag R2 gave it.
186#[derive(Deserialize)]
187struct Part {
188 part: u16,
189 etag: String,
190}
191
192#[derive(Deserialize)]
193struct Complete {
194 size: u64,
195 parts: Vec<Part>,
196}
197
198/// An outcome of the actions service, or its failure as the reply it means.
199fn refused<T>(outcome: Outcome<T>) -> std::result::Result<T, Result<Response>> {
200 match outcome {
201 Outcome::Ok(value) => Ok(value),
202 Outcome::Fail(failure) => {
203 let status = match failure.code {
204 FailureCode::Unauthenticated => 401,
205 FailureCode::NotFound => 404,
206 FailureCode::Conflict => 409,
207 FailureCode::Forbidden => 403,
208 _ => 400,
209 };
210 Err(error(status, &failure.message))
211 }
212 }
213}
214
215/// The cache: restoring, uploading in parts, and the older whole upload.
216#[allow(clippy::too_many_arguments)]
217async fn cache(
218 mut request: Request,
219 bucket: &Bucket,
220 services: &Services,
221 method: &str,
222 job: &str,
223 token: &str,
224 repo: &str,
225 what: &str,
226) -> Result<Response> {
227 let parts: Vec<&str> = what.split('/').collect();
228 let upload_id = query(&request, "upload").unwrap_or_default();
Merge branch 'worktree-agent-a3abfcce648e87dca'229 // The hash of the entry's paths and compression; runners from before
230 // it was sent send none.
231 let version = query(&request, "version").unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents232 match (method, parts.as_slice()) {
233 ("GET", ["cache"]) => {
234 let key = query(&request, "key").unwrap_or_default();
235 let restore: Vec<String> =
236 query(&request, "restore").unwrap_or_default().lines().map(str::trim).filter(|p| !p.is_empty()).map(str::to_owned).collect();
237 let found: Outcome<Option<CacheHit>> = g1t_kit::call(
238 &services.actions,
239 "cache_lookup",
Merge branch 'worktree-agent-a3abfcce648e87dca'240 &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 incidents241 )
242 .await?;
243 let found = match refused(found) {
244 Ok(found) => found,
245 Err(reply) => return reply,
246 };
247 if let Some(hit) = found
248 && let Some(object) = bucket.get(&hit.object).execute().await?
249 && let Some(body) = object.body()
250 {
251 let mut response = Response::from_body(body.response_body()?)?;
252 let headers = response.headers_mut();
253 headers.set("x-g1t-key", &hit.key)?;
254 headers.set("content-length", &object.size().to_string())?;
255 headers.set("content-type", "application/octet-stream")?;
256 return Ok(response);
257 }
Merge branch 'worktree-agent-a3abfcce648e87dca'258 // Entries kept in KV before the cache moved to R2 had no scope,
259 // so they are never restored.
260 error(404, "Nothing cached under those keys.")
Fast pages, required checks on the branch, self-hosted runners, honest incidents261 }
262 // Older runners send a whole entry of at most 60 MB at once.
263 ("PUT", ["cache"]) => {
A repository has its own sidebar, as settings do264 let key = query(&request, "key").unwrap_or_default();
Fast pages, required checks on the branch, self-hosted runners, honest incidents265 let bytes = request.bytes().await?;
266 if bytes.len() > MAX_BYTES {
267 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 do268 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents269 let reserved: Outcome<CacheReservation> = g1t_kit::call(
270 &services.actions,
271 "cache_reserve",
Merge branch 'worktree-agent-a3abfcce648e87dca'272 &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 incidents273 )
274 .await?;
275 let reserved = match reserved {
276 Outcome::Fail(failure) if failure.code == FailureCode::Conflict => {
277 return crate::reply(&json!({ "saved": false, "reason": failure.message }));
A repository has its own sidebar, as settings do278 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents279 other => match refused(other) {
280 Ok(reserved) => reserved,
281 Err(reply) => return reply,
282 },
283 };
284 let size = bytes.len() as u64;
285 bucket.put(&reserved.object, bytes).execute().await?;
286 commit(bucket, services, job, token, &reserved.id, size).await?;
287 crate::reply(&json!({ "saved": true }))
288 }
289 ("POST", ["cache", "uploads"]) => {
290 let key = query(&request, "key").unwrap_or_default();
291 let size = query(&request, "size").and_then(|s| s.parse::<u64>().ok()).unwrap_or(0);
292 let reserved: Outcome<CacheReservation> = g1t_kit::call(
293 &services.actions,
294 "cache_reserve",
Merge branch 'worktree-agent-a3abfcce648e87dca'295 &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 incidents296 )
297 .await?;
298 let reserved = match refused(reserved) {
299 Ok(reserved) => reserved,
300 Err(reply) => return reply,
301 };
302 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
303 crate::reply(&json!({ "id": reserved.id, "upload": upload.upload_id().await, "part_bytes": CACHE_PART_BYTES }))
304 }
305 ("PUT", ["cache", "uploads", id, part]) => {
306 let part = part.parse::<u16>().unwrap_or(0);
307 if part == 0 || upload_id.is_empty() {
308 return error(400, "A part is numbered from 1, and names its upload.");
A repository has its own sidebar, as settings do309 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents310 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
311 if length == 0 || length > CACHE_PART_BYTES {
312 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 do313 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents314 // The body goes to R2 as it comes, never held whole here.
315 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
316 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
317 let uploaded = upload.upload_part(part, body).await?;
318 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
319 }
320 ("POST", ["cache", "uploads", id, "complete"]) => {
321 let done: Complete = match request.json().await {
322 Ok(done) => done,
323 Err(_) => return error(400, "Send { size, parts: [{ part, etag }] }."),
324 };
325 let upload = bucket.resume_multipart_upload(object_of(repo, id), &upload_id)?;
326 let mut parts = done.parts;
327 parts.sort_by_key(|p| p.part);
328 if let Err(problem) = upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
329 let _ = abort(services, job, token, id).await;
330 return error(400, &format!("The upload could not be completed: {problem}"));
331 }
332 commit(bucket, services, job, token, id, done.size).await?;
333 crate::reply(&json!({ "saved": true }))
334 }
335 ("DELETE", ["cache", "uploads", id]) => {
336 if let Ok(upload) = bucket.resume_multipart_upload(object_of(repo, id), &upload_id) {
337 let _ = upload.abort().await;
A repository has its own sidebar, as settings do338 }
Fast pages, required checks on the branch, self-hosted runners, honest incidents339 abort(services, job, token, id).await?;
340 crate::reply(&json!({ "aborted": true }))
A repository has its own sidebar, as settings do341 }
342 _ => error(404, "No such endpoint."),
343 }
344}
345
Fast pages, required checks on the branch, self-hosted runners, honest incidents346/// Where an entry is in R2: under its repository, by its id, as the
347/// actions service named it when it was reserved.
348fn object_of(repo: &str, id: &str) -> String {
349 format!("c/{repo}/{id}")
350}
351
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2352/// Where an artifact is in R2, as the actions service names it.
353fn artifact_object(repo: &str, id: u64) -> String {
354 format!("a/{repo}/{id}")
355}
356
357/// An artifact as a job's runner lists it.
358fn for_runner(artifact: &Artifact) -> Value {
359 json!({
360 "id": artifact.id,
361 "name": artifact.name,
362 "size": artifact.size,
363 "digest": artifact.digest,
364 "format": artifact.format,
365 "created_at": artifact.created_at,
366 "expires_at": artifact.expires_at,
367 })
368}
369
370/// An R2 object streamed back, with what it is.
371async fn stream(bucket: &Bucket, object: &str, format: &str) -> Result<Response> {
372 let Some(found) = bucket.get(object).execute().await? else { return error(404, "That artifact is gone: it expired or was deleted.") };
373 let size = found.size();
374 let Some(body) = found.body() else { return error(404, "That artifact is gone: it expired or was deleted.") };
375 let mut response = Response::from_body(body.response_body()?)?;
376 let headers = response.headers_mut();
377 headers.set("content-length", &size.to_string())?;
378 headers.set("content-type", if format == "tgz" { "application/gzip" } else { "application/zip" })?;
379 headers.set("x-g1t-format", format)?;
380 Ok(response)
381}
382
383/// A job's artifacts: listing its run's (or another run's of its
384/// repository), uploading in parts, downloading, deleting; and the whole
385/// uploads and downloads by name of older runners.
386#[allow(clippy::too_many_arguments)]
387async fn artifacts(
388 mut request: Request,
389 kv: &KvStore,
390 bucket: &Bucket,
391 services: &Services,
392 method: &str,
393 job: &str,
394 token: &str,
395 run: &str,
396 repo: &str,
397 what: &str,
398) -> Result<Response> {
399 let parts: Vec<&str> = what.split('/').collect();
400 let upload_id = query(&request, "upload").unwrap_or_default();
401 let credential = |id: Option<u64>, name: Option<String>, run_id: Option<String>| JobArtifactsArgs {
402 job: job.to_owned(),
403 token: token.to_owned(),
404 run_id,
405 name,
406 id,
407 };
408 let actions = &services.actions;
409 match (method, parts.as_slice()) {
410 ("GET", ["artifacts"]) => {
411 let run_id = query(&request, "run_id").filter(|r| !r.is_empty());
412 let found: Outcome<Vec<Artifact>> = g1t_kit::call(actions, "job_artifacts", &credential(None, None, run_id.clone())).await?;
413 match refused(found) {
414 Ok(found) => {
415 let mut listed: Vec<Value> = found.iter().map(for_runner).collect();
416 // Artifacts older runners kept in KV, for the days they
417 // are still there.
418 let legacy_run = run_id.as_deref().unwrap_or(run);
Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails419 if legacy_kv(g1t_kit::now_ms()) {
420 for (_, meta) in list(kv, &format!("a/{legacy_run}/")).await? {
421 if !listed.iter().any(|a| a["name"] == meta.name.as_str()) {
422 listed.push(json!({ "name": meta.name, "size": meta.size, "format": "tgz" }));
423 }
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2424 }
425 }
426 crate::reply(&listed)
427 }
428 Err(reply) => reply,
429 }
430 }
431 ("POST", ["artifacts", "uploads"]) => {
432 let flag = |name: &str| query(&request, name).is_some_and(|v| v == "true");
433 let args = ArtifactReserveArgs {
434 job: job.to_owned(),
435 token: token.to_owned(),
436 name: query(&request, "name").unwrap_or_default(),
437 size: query(&request, "size").and_then(|s| s.parse().ok()).unwrap_or(0),
438 retention_days: query(&request, "retention_days").and_then(|d| d.parse().ok()).unwrap_or(0),
439 expires_at: None,
440 overwrite: flag("overwrite"),
441 format: query(&request, "format"),
442 };
443 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
444 let reserved = match refused(reserved) {
445 Ok(reserved) => reserved,
446 Err(reply) => return reply,
447 };
448 let upload = bucket.create_multipart_upload(&reserved.object).execute().await?;
449 crate::reply(&json!({
450 "id": reserved.id,
451 "upload": upload.upload_id().await,
452 "part_bytes": ARTIFACT_PART_BYTES,
453 "retention_days": reserved.retention_days,
454 "expires_at": reserved.expires_at,
455 }))
456 }
457 ("PUT", ["artifacts", "uploads", id, part]) => {
458 let (Ok(id), Ok(part)) = (id.parse::<u64>(), part.parse::<u16>()) else {
459 return error(400, "A part is numbered from 1, of an artifact named by its number.");
460 };
461 if part == 0 || upload_id.is_empty() {
462 return error(400, "A part is numbered from 1, and names its upload.");
463 }
464 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
465 if length == 0 || length > ARTIFACT_PART_BYTES {
466 return error(413, &format!("A part is 1 to {} MB, with its length.", ARTIFACT_PART_BYTES / 1_048_576));
467 }
468 let Some(body) = request.inner().body() else { return error(400, "The part is empty.") };
469 let upload = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id)?;
470 let uploaded = upload.upload_part(part, body).await?;
471 crate::reply(&json!({ "part": uploaded.part_number(), "etag": uploaded.etag() }))
472 }
473 ("POST", ["artifacts", "uploads", id, "complete"]) => {
474 let Ok(id) = id.parse::<u64>() else { return error(404, "No such upload.") };
475 let done: Value = request.json().await.unwrap_or(Value::Null);
476 let mut parts: Vec<Part> = serde_json::from_value(done["parts"].clone()).unwrap_or_default();
477 if parts.is_empty() {
478 return error(400, "Send { size, parts: [{ part, etag }], digest }.");
479 }
480 parts.sort_by_key(|p| p.part);
481 let upload = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id)?;
482 let object = match upload.complete(parts.into_iter().map(|p| UploadedPart::new(p.part, p.etag))).await {
483 Ok(object) => object,
484 Err(problem) => {
485 let _: Outcome<bool> = g1t_kit::call(actions, "artifact_abort", &credential(Some(id), None, None)).await?;
486 return error(400, &format!("The upload could not be completed: {problem}"));
487 }
488 };
489 let args = ArtifactCommitArgs {
490 job: job.to_owned(),
491 token: token.to_owned(),
492 id: Some(id),
493 name: None,
494 // What R2 holds, not what the runner says.
495 size: object.size(),
496 digest: done["digest"].as_str().map(str::to_owned),
497 };
498 let committed: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
499 match refused(committed) {
500 Ok(artifact) => crate::reply(&for_runner(&artifact)),
501 Err(reply) => reply,
502 }
503 }
504 ("DELETE", ["artifacts", "uploads", id]) => {
505 let Ok(id) = id.parse::<u64>() else { return error(404, "No such upload.") };
506 if let Ok(upload) = bucket.resume_multipart_upload(artifact_object(repo, id), &upload_id) {
507 let _ = upload.abort().await;
508 }
509 let _: Outcome<bool> = g1t_kit::call(actions, "artifact_abort", &credential(Some(id), None, None)).await?;
510 crate::reply(&json!({ "aborted": true }))
511 }
512 ("GET", ["artifacts", id, "download"]) => {
513 let Ok(id) = id.parse::<u64>() else { return error(404, "No such artifact.") };
514 // Any run of the job's repository: the service checks.
515 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &credential(Some(id), None, query(&request, "run_id"))).await?;
516 match found {
517 Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await,
518 Outcome::Fail(failure) => error(failure.code.http_status(), &failure.message),
519 }
520 }
521 ("DELETE", ["artifacts", id]) => {
522 let Ok(id) = id.parse::<u64>() else { return error(404, "No such artifact.") };
523 let done: Outcome<Artifact> = g1t_kit::call(actions, "job_delete_artifact", &credential(Some(id), None, None)).await?;
524 match refused(done) {
525 Ok(artifact) => crate::reply(&for_runner(&artifact)),
526 Err(reply) => reply,
527 }
528 }
529 // Older runners: a whole artifact of at most 60 MB, by name.
530 ("PUT", ["artifacts", name]) => {
531 let name = decode(name);
532 let bytes = request.bytes().await?;
533 if bytes.len() > MAX_BYTES {
534 return error(413, "An artifact sent at once is at most 60 MB; newer runners upload it in parts.");
535 }
536 let args = ArtifactReserveArgs {
537 job: job.to_owned(),
538 token: token.to_owned(),
539 name: name.clone(),
540 size: bytes.len() as u64,
541 format: Some("tgz".to_owned()),
542 ..ArtifactReserveArgs::default()
543 };
544 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
545 let reserved = match refused(reserved) {
546 Ok(reserved) => reserved,
547 Err(reply) => return reply,
548 };
549 let size = bytes.len() as u64;
550 bucket.put(&reserved.object, bytes).execute().await?;
551 let args = ArtifactCommitArgs { job: job.to_owned(), token: token.to_owned(), id: Some(reserved.id), name: None, size, digest: None };
552 let committed: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
553 match refused(committed) {
554 Ok(_) => crate::reply(&json!({ "name": name, "size": size })),
555 Err(reply) => reply,
556 }
557 }
558 ("GET", ["artifacts", name]) => {
559 let name = decode(name);
560 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &credential(None, Some(name.clone()), None)).await?;
561 match found {
562 Outcome::Ok(found) => stream(bucket, &found.object, &found.artifact.format).await,
563 // One an older runner kept in KV.
Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails564 Outcome::Fail(_) if legacy_kv(g1t_kit::now_ms()) => match get(kv, &format!("a/{run}/{name}")).await? {
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2565 Some(bytes) => Response::from_bytes(bytes),
566 None => error(404, "No such artifact."),
567 },
Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails568 Outcome::Fail(_) => error(404, "No such artifact."),
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2569 }
570 }
571 _ => error(404, "No such endpoint."),
572 }
573}
574
Fast pages, required checks on the branch, self-hosted runners, honest incidents575/// Marks an uploaded entry ready, and deletes what that evicted.
576async fn commit(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<()> {
577 let committed: Outcome<CacheCommitted> = g1t_kit::call(
578 &services.actions,
579 "cache_commit",
580 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
581 )
582 .await?;
583 if let Outcome::Ok(committed) = committed
584 && !committed.evicted.is_empty()
585 {
586 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
587 }
588 Ok(())
589}
590
591async fn abort(services: &Services, job: &str, token: &str, id: &str) -> Result<()> {
592 let _: Outcome<bool> = g1t_kit::call(
593 &services.actions,
594 "cache_abort",
595 &CacheAbortArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned() },
596 )
597 .await?;
598 Ok(())
599}
600
601/// An entry saved in KV before the cache moved to R2, by key or restore key.
A repository has its own sidebar, as settings do602/// Someone who can see the run downloading one of its artifacts.
603pub async fn download(env: &Env, services: &Services, viewer: &g1t_contracts::Viewer, owner: &str, repo: &str, run: &str, name: &str) -> Result<Response> {
604 let seen: Outcome<Value> = g1t_kit::call(
605 &services.actions,
606 "run",
607 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
608 )
609 .await?;
610 if let Outcome::Fail(refused) = seen {
611 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
612 return error(status, &refused.message);
613 }
614 let name = decode(name);
Actions: OIDC tokens, the toolkit's cache and artifact services, and artifacts in R2615 // Kept in R2: a redirect to a link signed for a few minutes.
616 let args = ArtifactArgs {
617 repo: g1t_contracts::repos::RepoPath { namespace: owner.to_owned(), name: repo.to_owned() },
618 viewer: viewer.clone(),
619 id: None,
620 run: Some(run.to_owned()),
621 name: Some(name.clone()),
622 };
623 let found: Outcome<ArtifactBlob> = g1t_kit::call(&services.actions, "artifact_download", &args).await?;
624 if let Outcome::Ok(found) = found
625 && !found.blob.is_empty()
626 {
627 return Response::redirect_with_status(Url::parse(&crate::artifacts::blob_url(&services.addresses.api, &found.blob))?, 302);
628 }
629 // Kept in KV by an older runner.
Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails630 if !legacy_kv(g1t_kit::now_ms()) {
631 return error(404, "No such artifact, or it has expired.");
632 }
A repository has its own sidebar, as settings do633 match get(&store(env)?, &format!("a/{run}/{name}")).await? {
634 Some(bytes) => {
635 let mut response = Response::from_bytes(bytes)?;
636 let headers = response.headers_mut();
637 headers.set("content-type", "application/gzip")?;
638 headers.set("content-disposition", &format!("attachment; filename=\"{}.tar.gz\"", name.replace('"', "")))?;
639 Ok(response)
640 }
641 None => error(404, "No such artifact, or it has expired."),
642 }
643}
644
Merge branch 'worktree-agent-ab9543c492a7ed481' into spend-guardrails645#[cfg(test)]
646mod tests {
647 use super::*;
648
649 #[test]
650 fn kv_artifacts_are_not_asked_for_after_they_have_all_expired() {
651 assert_eq!(g1t_contracts::time::rfc3339(LEGACY_KV_UNTIL_MS), "2026-10-22T00:00:00.000Z");
652 assert!(legacy_kv(LEGACY_KV_UNTIL_MS - 1));
653 assert!(!legacy_kv(LEGACY_KV_UNTIL_MS));
654 }
655}

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