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

Merge main into the run-protection branch1//! The services GitHub's Actions toolkit calls from inside a job, so that
2//! actions built on `@actions/cache` and `@actions/artifact` (such as
3//! `actions/setup-node` with `cache:`, or `Swatinem/rust-cache`) work on
4//! g1t unchanged. A job is told where they are in its variables
5//! (`ACTIONS_RUNTIME_TOKEN`, `ACTIONS_RESULTS_URL`, `ACTIONS_CACHE_URL`,
6//! `ACTIONS_CACHE_SERVICE_V2`; see `runtime_variables`).
7//!
8//! - Twirp, at `/twirp/github.actions.results.api.v1.CacheService/…` and
9//! `…ArtifactService/…`: the cache's newer protocol (`CreateCacheEntry`,
10//! `FinalizeCacheEntryUpload`, `GetCacheEntryDownloadURL`) and the
11//! artifacts' (`CreateArtifact`, `FinalizeArtifact`, `ListArtifacts`,
12//! `GetSignedArtifactURL`, `DeleteArtifact`). JSON, the toolkit's field
13//! names.
14//! - The cache's older protocol, at `{ACTIONS_CACHE_URL}_apis/artifactcache/…`,
15//! which the toolkit's client uses whenever the server it runs against is
16//! not github.com: on g1t, that is the one it uses.
17//! - Blobs, at `/actions/toolkit/blobs/{token}`: the signed links those
18//! hand out. Downloads are a plain GET. Uploads speak the part of Azure
19//! Blob Storage's protocol the toolkit's client uses (Put Blob, Put
20//! Block, Put Block List), mapped onto an R2 multipart upload: a block's
21//! id ends in its index, which is its part's number.
22//!
23//! Every call carries the job's runtime token; the actions service checks
24//! it and keeps the entries (cache.rs, artifacts.rs, runtime.rs there).
25
26use base64::Engine;
27use base64::engine::general_purpose::{STANDARD, URL_SAFE_NO_PAD};
28use g1t_contracts::actions::{
29 ARTIFACT_MAX_BYTES, Artifact, ArtifactBlob, ArtifactCommitArgs, ArtifactReservation, ArtifactReserveArgs, BlobArgs, BlobGrant, BlobPart,
30 BlobSignArgs, CACHE_MAX_ENTRY_BYTES, CACHE_PART_BYTES, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation,
31 CacheReserveArgs, CacheUploadArgs, JobArtifactsArgs,
32};
33use g1t_contracts::{Failure, FailureCode, Outcome};
34use serde_json::{Map, Value, json};
35use worker::{Bucket, Env, Request, Response, Result, UploadedPart};
36
37use crate::artifacts::blob_url;
38use crate::operations::Services;
39
40/// The Twirp services, by name.
41pub const CACHE_SERVICE: &str = "github.actions.results.api.v1.CacheService";
42pub const ARTIFACT_SERVICE: &str = "github.actions.results.api.v1.ArtifactService";
43/// Where the cache's older protocol is: `ACTIONS_CACHE_URL`.
44pub const CACHE_PATH: &str = "/actions/toolkit/";
45/// The largest single block or blob a request may carry.
46const MAX_BLOCK_BYTES: u64 = 256 * 1024 * 1024;
47
48/// The bearer token of a request.
49pub fn bearer(request: &Request) -> String {
50 request
51 .headers()
52 .get("authorization")
53 .ok()
54 .flatten()
55 .and_then(|h| h.split_once(' ').map(|(_, t)| t.trim().to_owned()))
56 .unwrap_or_default()
57}
58
59fn claims(token: &str) -> Option<Value> {
60 serde_json::from_slice(&URL_SAFE_NO_PAD.decode(token.split('.').nth(1)?).ok()?).ok()
61}
62
63/// The job a runtime token names, unchecked: the actions service checks it.
64pub fn runtime_job(token: &str) -> Option<String> {
65 claims(token)?["job"].as_str().map(str::to_owned)
66}
67
68/// The variables a job gets for the toolkit: its runtime token and where
69/// the services are, and where to ask for an OIDC token when it may.
70pub fn runtime_variables(api: &str, token: &str, id_token: bool) -> Map<String, Value> {
71 let mut vars = Map::new();
72 let mut set = |k: &str, v: String| {
73 vars.insert(k.to_owned(), Value::String(v));
74 };
75 set("ACTIONS_RUNTIME_TOKEN", token.to_owned());
76 set("ACTIONS_RESULTS_URL", format!("{api}/"));
77 set("ACTIONS_CACHE_URL", format!("{api}{CACHE_PATH}"));
78 set("ACTIONS_CACHE_SERVICE_V2", "True".to_owned());
79 if id_token {
80 set("ACTIONS_ID_TOKEN_REQUEST_URL", format!("{}/token?api-version=2.0", crate::oidc::issuer(api)));
81 set("ACTIONS_ID_TOKEN_REQUEST_TOKEN", token.to_owned());
82 }
83 vars
84}
85
86// ── Twirp ───────────────────────────────────────────────────────────────────
87
88/// A Twirp error: its code and message, at the status Twirp gives it.
89fn twirp_error(code: &str, message: &str) -> Result<Response> {
90 let status = match code {
91 "unauthenticated" => 401,
92 "permission_denied" => 403,
93 "not_found" => 404,
94 "already_exists" => 409,
95 "invalid_argument" => 400,
96 "failed_precondition" => 412,
97 "resource_exhausted" => 429,
98 _ => 500,
99 };
100 Ok(Response::from_json(&json!({ "code": code, "msg": message }))?.with_status(status))
101}
102
103fn twirp_failure(failure: &Failure) -> Result<Response> {
104 let code = match failure.code {
105 FailureCode::Unauthenticated => "unauthenticated",
106 FailureCode::Forbidden => "permission_denied",
107 FailureCode::NotFound => "not_found",
108 FailureCode::Conflict => "already_exists",
109 FailureCode::Invalid => "invalid_argument",
110 _ => "failed_precondition",
111 };
112 twirp_error(code, &failure.message)
113}
114
115/// A field in the toolkit's spelling (`snake_case`), or its JSON name.
116fn field<'a>(body: &'a Value, name: &str) -> &'a Value {
117 if !body[name].is_null() {
118 return &body[name];
119 }
120 let camel: String = name.split('_').enumerate().map(|(i, p)| if i == 0 { p.to_owned() } else { p[..1].to_uppercase() + &p[1..] }).collect();
121 &body[camel]
122}
123
124fn text(body: &Value, name: &str) -> String {
125 match field(body, name) {
126 Value::String(s) => s.clone(),
127 Value::Number(n) => n.to_string(),
128 // A wrapper written as an object, `{ "value": … }`.
129 Value::Object(o) => o.get("value").map(|v| v.as_str().map_or_else(|| v.to_string(), str::to_owned)).unwrap_or_default(),
130 _ => String::new(),
131 }
132}
133
134fn number(body: &Value, name: &str) -> Option<u64> {
135 text(body, name).trim().parse().ok()
136}
137
138/// The run and job a runtime token names, which a request's backend ids
139/// must match.
140fn backend_ids(token: &str) -> (String, String) {
141 let c = claims(token).unwrap_or_default();
142 (c["run"].as_str().unwrap_or_default().to_owned(), c["job"].as_str().unwrap_or_default().to_owned())
143}
144
145/// What the toolkit's artifact client lists.
146fn listed(artifact: &Artifact) -> Value {
147 json!({
148 "workflow_run_backend_id": artifact.run_id,
149 "workflow_job_run_backend_id": artifact.job_id,
150 "database_id": artifact.id.to_string(),
151 "name": artifact.name,
152 "size": artifact.size.to_string(),
153 "created_at": artifact.created_at,
154 "digest": artifact.digest,
155 })
156}
157
158/// Starts an R2 upload for an entry the service reserved, and the signed
159/// link the toolkit sends it to.
160async fn start_upload(bucket: &Bucket, services: &Services, job: &str, token: &str, kind: &str, id: &str, object: &str) -> Result<Outcome<String>> {
161 let upload = bucket.create_multipart_upload(object).execute().await?;
162 let upload = upload.upload_id().await;
163 let signed: Outcome<String> = g1t_kit::call(
164 &services.actions,
165 "blob_sign",
166 &BlobSignArgs { job: job.to_owned(), token: token.to_owned(), kind: kind.to_owned(), id: id.to_owned(), upload },
167 )
168 .await?;
169 Ok(match signed {
170 Outcome::Ok(blob) => Outcome::Ok(blob_url(&services.addresses.api, &blob)),
171 Outcome::Fail(refused) => Outcome::Fail(refused),
172 })
173}
174
175/// `POST /twirp/{service}/{method}`.
176pub async fn twirp(mut request: Request, env: &Env, services: &Services, service: &str, method: &str) -> Result<Response> {
177 let token = bearer(&request);
178 let Some(job) = runtime_job(&token) else {
179 return twirp_error("unauthenticated", "Send the job's ACTIONS_RUNTIME_TOKEN as a bearer token.");
180 };
181 let body: Value = request.json().await.unwrap_or(Value::Null);
182 let bucket = env.bucket("ACTIONS_CACHE")?;
183 let actions = &services.actions;
184 let (run, own_job) = backend_ids(&token);
185 // An artifact call names its run, and for an upload its job: the
186 // token's.
187 if service == ARTIFACT_SERVICE {
188 let asked_run = text(&body, "workflow_run_backend_id");
189 if !asked_run.is_empty() && asked_run != run {
190 return twirp_error("permission_denied", "The runtime token is for another run.");
191 }
192 let asked_job = text(&body, "workflow_job_run_backend_id");
193 if matches!(method, "CreateArtifact" | "FinalizeArtifact") && !asked_job.is_empty() && asked_job != own_job {
194 return twirp_error("permission_denied", "The runtime token is for another job.");
195 }
196 }
197 match (service, method) {
198 (CACHE_SERVICE, "GetCacheEntryDownloadURL") => {
199 let restore: Vec<String> = field(&body, "restore_keys").as_array().map(|k| k.iter().filter_map(|v| v.as_str().map(str::to_owned)).collect()).unwrap_or_default();
200 let args = CacheLookupArgs { job, token, key: text(&body, "key"), restore, version: Some(text(&body, "version")) };
201 let found: Outcome<Option<CacheHit>> = g1t_kit::call(actions, "cache_lookup", &args).await?;
202 match found {
203 Outcome::Ok(Some(CacheHit { key, blob: Some(blob), .. })) => {
204 Response::from_json(&json!({ "ok": true, "signed_download_url": blob_url(&services.addresses.api, &blob), "matched_key": key }))
205 }
206 Outcome::Ok(_) => Response::from_json(&json!({ "ok": false, "signed_download_url": "", "matched_key": "" })),
207 Outcome::Fail(refused) => twirp_failure(&refused),
208 }
209 }
210 (CACHE_SERVICE, "CreateCacheEntry") => {
211 let args = CacheReserveArgs { job: job.clone(), token: token.clone(), key: text(&body, "key"), size: 0, version: Some(text(&body, "version")) };
212 let reserved: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_reserve", &args).await?;
213 let reserved = match reserved {
214 Outcome::Ok(reserved) => reserved,
215 // The client warns with this and goes on, as for a key
216 // another job is saving.
217 Outcome::Fail(refused) => return Response::from_json(&json!({ "ok": false, "signed_upload_url": "", "message": refused.message })),
218 };
219 match start_upload(&bucket, services, &job, &token, "cache", &reserved.id, &reserved.object).await? {
220 Outcome::Ok(url) => Response::from_json(&json!({ "ok": true, "signed_upload_url": url })),
221 Outcome::Fail(refused) => Response::from_json(&json!({ "ok": false, "signed_upload_url": "", "message": refused.message })),
222 }
223 }
224 (CACHE_SERVICE, "FinalizeCacheEntryUpload") => {
225 let args = CacheUploadArgs { job: job.clone(), token: token.clone(), number: None, key: Some(text(&body, "key")), version: Some(text(&body, "version")) };
226 let pending: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_upload", &args).await?;
227 let pending = match pending {
228 Outcome::Ok(pending) => pending,
229 Outcome::Fail(refused) => return Response::from_json(&json!({ "ok": false, "entry_id": "0", "message": refused.message })),
230 };
231 let Some(object) = bucket.head(&pending.object).await? else {
232 return Response::from_json(&json!({ "ok": false, "entry_id": "0", "message": "Nothing was uploaded for that entry." }));
233 };
234 match commit_cache(&bucket, services, &job, &token, &pending.id, object.size()).await? {
235 Outcome::Ok(()) => Response::from_json(&json!({ "ok": true, "entry_id": pending.number.to_string() })),
236 Outcome::Fail(refused) => Response::from_json(&json!({ "ok": false, "entry_id": "0", "message": refused.message })),
237 }
238 }
239 (ARTIFACT_SERVICE, "CreateArtifact") => {
240 let expires_at = Some(text(&body, "expires_at")).filter(|e| !e.is_empty());
241 let args = ArtifactReserveArgs { job: job.clone(), token: token.clone(), name: text(&body, "name"), expires_at, ..ArtifactReserveArgs::default() };
242 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
243 let reserved = match reserved {
244 Outcome::Ok(reserved) => reserved,
245 Outcome::Fail(refused) => return twirp_failure(&refused),
246 };
247 match start_upload(&bucket, services, &job, &token, "artifact", &reserved.id.to_string(), &reserved.object).await? {
248 Outcome::Ok(url) => Response::from_json(&json!({ "ok": true, "signed_upload_url": url })),
249 Outcome::Fail(refused) => twirp_failure(&refused),
250 }
251 }
252 (ARTIFACT_SERVICE, "FinalizeArtifact") => {
253 let digest = Some(text(&body, "hash")).filter(|h| !h.is_empty());
254 // The size it was measured at as it was stored, not the one it
255 // says.
256 let args = ArtifactCommitArgs { job, token, id: None, name: Some(text(&body, "name")), size: 0, digest };
257 let done: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
258 match done {
259 Outcome::Ok(artifact) => Response::from_json(&json!({ "ok": true, "artifact_id": artifact.id.to_string() })),
260 Outcome::Fail(refused) => twirp_failure(&refused),
261 }
262 }
263 (ARTIFACT_SERVICE, "ListArtifacts") => {
264 let name = Some(text(&body, "name_filter")).filter(|n| !n.is_empty());
265 let id = number(&body, "id_filter");
266 let args = JobArtifactsArgs { job, token, run_id: None, name, id };
267 let found: Outcome<Vec<Artifact>> = g1t_kit::call(actions, "job_artifacts", &args).await?;
268 match found {
269 Outcome::Ok(found) => Response::from_json(&json!({ "artifacts": found.iter().map(listed).collect::<Vec<_>>() })),
270 Outcome::Fail(refused) => twirp_failure(&refused),
271 }
272 }
273 (ARTIFACT_SERVICE, "GetSignedArtifactURL") => {
274 let args = JobArtifactsArgs { job, token, run_id: None, name: Some(text(&body, "name")), id: None };
275 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &args).await?;
276 match found {
277 Outcome::Ok(found) if !found.blob.is_empty() => Response::from_json(&json!({ "signed_url": blob_url(&services.addresses.api, &found.blob) })),
278 Outcome::Ok(_) => twirp_error("failed_precondition", "Download links are not set up on this installation."),
279 Outcome::Fail(refused) => twirp_failure(&refused),
280 }
281 }
282 (ARTIFACT_SERVICE, "DeleteArtifact") => {
283 let args = JobArtifactsArgs { job, token, run_id: None, name: Some(text(&body, "name")), id: None };
284 let done: Outcome<Artifact> = g1t_kit::call(actions, "job_delete_artifact", &args).await?;
285 match done {
286 Outcome::Ok(artifact) => Response::from_json(&json!({ "ok": true, "artifact_id": artifact.id.to_string() })),
287 Outcome::Fail(refused) => twirp_failure(&refused),
288 }
289 }
290 _ => twirp_error("bad_route", &format!("No method {method} on {service}.")),
291 }
292}
293
294/// Marks an uploaded cache entry ready, and deletes what that evicted.
295async fn commit_cache(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<Outcome<()>> {
296 let committed: Outcome<CacheCommitted> = g1t_kit::call(
297 &services.actions,
298 "cache_commit",
299 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
300 )
301 .await?;
302 Ok(match committed {
303 Outcome::Ok(committed) => {
304 if !committed.evicted.is_empty() {
305 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
306 }
307 Outcome::Ok(())
308 }
309 Outcome::Fail(refused) => Outcome::Fail(refused),
310 })
311}
312
313// ── The cache's older protocol ──────────────────────────────────────────────
314
315fn plain_error(status: u16, message: &str) -> Result<Response> {
316 Ok(Response::from_json(&json!({ "message": message, "error": { "message": message } }))?.with_status(status))
317}
318
319fn query(request: &Request, name: &str) -> Option<String> {
320 request.url().ok()?.query_pairs().find(|(k, _)| k == name).map(|(_, v)| v.into_owned())
321}
322
323/// The part a chunk of the older protocol is, from its `Content-Range`:
324/// chunks are `CACHE_PART_BYTES` apart, as the toolkit sends them.
325pub fn chunk_part(range: &str) -> Option<(u16, u64)> {
326 let range = range.trim().strip_prefix("bytes ")?;
327 let (span, _) = range.split_once('/')?;
328 let (start, end) = span.split_once('-')?;
329 let (start, end): (u64, u64) = (start.trim().parse().ok()?, end.trim().parse().ok()?);
330 if end < start || start % CACHE_PART_BYTES != 0 || end - start + 1 > CACHE_PART_BYTES {
331 return None;
332 }
333 Some(((start / CACHE_PART_BYTES + 1) as u16, end - start + 1))
334}
335
336/// `{ACTIONS_CACHE_URL}_apis/artifactcache/…`. `rest` is the path after it.
337pub async fn cache_v1(mut request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
338 let token = bearer(&request);
339 let Some(job) = runtime_job(&token) else {
340 return plain_error(401, "Send the job's ACTIONS_RUNTIME_TOKEN as a bearer token.");
341 };
342 let bucket = env.bucket("ACTIONS_CACHE")?;
343 let actions = &services.actions;
344 let parts: Vec<&str> = rest.split('/').filter(|p| !p.is_empty()).collect();
345 match (method, parts.as_slice()) {
346 ("GET", ["cache"]) => {
347 let keys: Vec<String> = query(&request, "keys").unwrap_or_default().split(',').map(|k| k.trim().to_owned()).filter(|k| !k.is_empty()).collect();
348 let Some((key, restore)) = keys.split_first() else {
349 return plain_error(400, "Give keys.");
350 };
351 let version = query(&request, "version").unwrap_or_default();
352 let args = CacheLookupArgs { job, token, key: key.clone(), restore: restore.to_vec(), version: Some(version.clone()) };
353 let found: Outcome<Option<CacheHit>> = g1t_kit::call(actions, "cache_lookup", &args).await?;
354 match found {
355 Outcome::Ok(Some(CacheHit { key, blob: Some(blob), created_at, .. })) => Response::from_json(&json!({
356 "cacheKey": key,
357 "cacheVersion": version,
358 "scope": "",
359 "creationTime": created_at,
360 "archiveLocation": blob_url(&services.addresses.api, &blob),
361 })),
362 Outcome::Ok(_) => Ok(Response::empty()?.with_status(204)),
363 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
364 }
365 }
366 ("POST", ["caches"]) => {
367 let body: Value = request.json().await.unwrap_or(Value::Null);
368 let size = body["cacheSize"].as_u64().unwrap_or(0);
369 let key = body["key"].as_str().unwrap_or_default().to_owned();
370 let version = body["version"].as_str().unwrap_or_default().to_owned();
371 let args = CacheReserveArgs { job: job.clone(), token: token.clone(), key, size, version: Some(version) };
372 let reserved: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_reserve", &args).await?;
373 let reserved = match reserved {
374 Outcome::Ok(reserved) => reserved,
375 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
376 };
377 match start_upload(&bucket, services, &job, &token, "cache", &reserved.id, &reserved.object).await? {
378 Outcome::Ok(_) => Ok(Response::from_json(&json!({ "cacheId": reserved.number }))?.with_status(201)),
379 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
380 }
381 }
382 (_, ["caches", number]) => {
383 let Ok(number) = number.parse::<u64>() else {
384 return plain_error(404, "No such cache entry.");
385 };
386 let args = CacheUploadArgs { job: job.clone(), token: token.clone(), number: Some(number), key: None, version: None };
387 let pending: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_upload", &args).await?;
388 let pending = match pending {
389 Outcome::Ok(pending) => pending,
390 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
391 };
392 let (Some(blob), Some(upload_id)) = (pending.blob.clone(), pending.upload.clone()) else {
393 return plain_error(409, "That entry's upload was not started.");
394 };
395 match method {
396 "PATCH" => {
397 let range = request.headers().get("content-range")?.unwrap_or_default();
398 let Some((part, length)) = chunk_part(&range) else {
399 return plain_error(400, &format!("Send the entry in chunks of {} MB, each with its Content-Range.", CACHE_PART_BYTES / 1_048_576));
400 };
401 let Some(body) = request.inner().body() else { return plain_error(400, "The chunk is empty.") };
402 let upload = bucket.resume_multipart_upload(&pending.object, &upload_id)?;
403 let uploaded = upload.upload_part(part, body).await?;
404 let recorded = BlobArgs { blob, part: u32::from(part), etag: uploaded.etag(), size: length };
405 let _: Outcome<bool> = g1t_kit::call(actions, "blob_part", &recorded).await?;
406 Ok(Response::empty()?.with_status(204))
407 }
408 "POST" => {
409 let parts: Outcome<Vec<BlobPart>> = g1t_kit::call(actions, "blob_parts", &BlobArgs { blob: blob.clone(), ..BlobArgs::default() }).await?;
410 let parts = match parts {
411 Outcome::Ok(parts) => parts,
412 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
413 };
414 let size = match finish(&bucket, &pending.object, &upload_id, &parts).await {
415 Ok(size) => size,
416 Err(problem) => return plain_error(400, &format!("The entry could not be completed: {problem}")),
417 };
418 let _: Outcome<bool> = g1t_kit::call(actions, "blob_done", &BlobArgs { blob, size, ..BlobArgs::default() }).await?;
419 match commit_cache(&bucket, services, &job, &token, &pending.id, size).await? {
420 Outcome::Ok(()) => Ok(Response::empty()?.with_status(204)),
421 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
422 }
423 }
424 _ => plain_error(405, "PATCH a chunk, or POST to commit."),
425 }
426 }
427 _ => plain_error(404, "No such endpoint."),
428 }
429}
430
431/// Completes an R2 upload from its recorded parts, in order; the size it
432/// came to. An upload with no parts is an empty object.
433async fn finish(bucket: &Bucket, object: &str, upload_id: &str, parts: &[BlobPart]) -> std::result::Result<u64, String> {
434 let upload = bucket.resume_multipart_upload(object, upload_id).map_err(|e| e.to_string())?;
435 if parts.is_empty() {
436 let _ = upload.abort().await;
437 bucket.put(object, Vec::<u8>::new()).execute().await.map_err(|e| e.to_string())?;
438 return Ok(0);
439 }
440 let done = upload
441 .complete(parts.iter().map(|p| UploadedPart::new(p.part as u16, p.etag.clone())))
442 .await
443 .map_err(|e| e.to_string())?;
444 Ok(done.size())
445}
446
447// ── Blobs ───────────────────────────────────────────────────────────────────
448
449/// A block's index, from its id: the toolkit's client (Azure's SDK) makes
450/// block ids as base64 of a prefix and the index padded with zeros.
451pub fn block_index(id: &str) -> Option<u32> {
452 let decoded = STANDARD.decode(id.trim()).ok()?;
453 let text = String::from_utf8(decoded).ok()?;
454 let digits: String = text.chars().rev().take_while(char::is_ascii_digit).collect::<Vec<_>>().into_iter().rev().collect();
455 if digits.is_empty() {
456 return None;
457 }
458 digits.parse().ok()
459}
460
461/// The block ids of a Put Block List body, in order.
462pub fn block_list(xml: &str) -> Vec<String> {
463 let mut ids = Vec::new();
464 let mut rest = xml;
465 while let Some(open) = rest.find('<') {
466 rest = &rest[open + 1..];
467 let Some(close) = rest.find('>') else { break };
468 let tag = &rest[..close];
469 rest = &rest[close + 1..];
470 if matches!(tag, "Latest" | "Committed" | "Uncommitted") {
471 let Some(end) = rest.find("</") else { break };
472 ids.push(rest[..end].trim().to_owned());
473 rest = &rest[end..];
474 }
475 }
476 ids
477}
478
479/// The parts a block list names, as recorded: each block must have been
480/// sent, as the part its index says, and they must run from the first.
481pub fn parts_for(ids: &[String], recorded: &[BlobPart]) -> std::result::Result<Vec<BlobPart>, String> {
482 let mut out = Vec::with_capacity(ids.len());
483 for (position, id) in ids.iter().enumerate() {
484 let index = block_index(id).ok_or_else(|| format!("The block id {id} does not end in its index."))?;
485 let part = index + 1;
486 if part as usize != position + 1 {
487 return Err("The blocks must be listed in the order they were numbered.".to_owned());
488 }
489 let found = recorded.iter().find(|p| p.part == part).ok_or_else(|| format!("Block {id} was never sent."))?;
490 out.push(found.clone());
491 }
492 Ok(out)
493}
494
495fn azure(status: u16) -> Result<Response> {
496 let mut response = Response::empty()?.with_status(status);
497 let headers = response.headers_mut();
498 headers.set("x-ms-request-id", &g1t_contracts::new_id("req", g1t_kit::now_ms()))?;
499 headers.set("x-ms-version", "2024-11-04")?;
500 headers.set("x-ms-request-server-encrypted", "true")?;
501 Ok(response)
502}
503
504fn azure_error(status: u16, code: &str, message: &str) -> Result<Response> {
505 let body = format!("<?xml version=\"1.0\" encoding=\"utf-8\"?><Error><Code>{code}</Code><Message>{message}</Message></Error>");
506 let mut response = Response::ok(body)?.with_status(status);
507 response.headers_mut().set("content-type", "application/xml")?;
508 response.headers_mut().set("x-ms-error-code", code)?;
509 Ok(response)
510}
511
512/// `/actions/toolkit/blobs/{token}`: GET or HEAD a download, PUT an upload.
513pub async fn blob(mut request: Request, env: &Env, services: &Services, method: &str, token: &str) -> Result<Response> {
514 let opened: Outcome<BlobGrant> = g1t_kit::call(&services.actions, "blob_open", &BlobArgs { blob: token.to_owned(), ..BlobArgs::default() }).await?;
515 let grant = match opened {
516 Outcome::Ok(grant) => grant,
517 Outcome::Fail(refused) => {
518 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
519 return azure_error(status, if status == 404 { "BlobNotFound" } else { "AuthenticationFailed" }, &refused.message);
520 }
521 };
522 let bucket = env.bucket("ACTIONS_CACHE")?;
523 match (method, grant.upload.as_deref()) {
524 ("GET" | "HEAD", None) => {
525 let headers = |response: &mut Response, size: u64| -> Result<()> {
526 let headers = response.headers_mut();
527 headers.set("content-length", &size.to_string())?;
528 headers.set("content-type", grant.content_type.as_deref().unwrap_or("application/octet-stream"))?;
529 headers.set("x-ms-blob-type", "BlockBlob")?;
530 if let Some(name) = &grant.filename {
531 headers.set("content-disposition", &format!("attachment; filename=\"{}\"", name.replace('"', "")))?;
532 }
533 Ok(())
534 };
535 if method == "HEAD" {
536 let Some(object) = bucket.head(&grant.object).await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
537 let mut response = Response::empty()?;
538 headers(&mut response, object.size())?;
539 return Ok(response);
540 }
541 let Some(object) = bucket.get(&grant.object).execute().await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
542 let size = object.size();
543 let Some(body) = object.body() else { return azure_error(404, "BlobNotFound", "It is gone.") };
544 let mut response = Response::from_body(body.response_body()?)?;
545 headers(&mut response, size)?;
546 Ok(response)
547 }
548 ("PUT", Some(upload_id)) => {
549 let comp = query(&request, "comp").unwrap_or_default();
550 let limit = if grant.kind == "cache" { CACHE_MAX_ENTRY_BYTES } else { ARTIFACT_MAX_BYTES };
551 match comp.as_str() {
552 // Put Block, or Put Blob: one part.
553 "block" | "" => {
554 let part = if comp == "block" {
555 match query(&request, "blockid").as_deref().and_then(block_index) {
556 Some(index) if index < 10_000 => index + 1,
557 _ => return azure_error(400, "InvalidQueryParameterValue", "A block id ends in its index, from 0."),
558 }
559 } else {
560 1
561 };
562 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
563 if length > MAX_BLOCK_BYTES || length > limit {
564 return azure_error(413, "RequestBodyTooLarge", "That block is larger than g1t takes at once.");
565 }
566 let upload = bucket.resume_multipart_upload(&grant.object, upload_id)?;
567 let uploaded = match request.inner().body() {
568 Some(body) if length > 0 => upload.upload_part(part as u16, body).await?,
569 _ => upload.upload_part(part as u16, Vec::<u8>::new()).await?,
570 };
571 let recorded = BlobArgs { blob: token.to_owned(), part, etag: uploaded.etag(), size: length };
572 let _: Outcome<bool> = g1t_kit::call(&services.actions, "blob_part", &recorded).await?;
573 if comp == "block" {
574 return azure(201);
575 }
576 // Put Blob is the whole thing: finish it now.
577 complete_blob(&bucket, services, token, &grant, upload_id, &[], limit, true).await
578 }
579 "blocklist" => {
580 let xml = request.text().await.unwrap_or_default();
581 let ids = block_list(&xml);
582 complete_blob(&bucket, services, token, &grant, upload_id, &ids, limit, false).await
583 }
584 _ => azure_error(400, "InvalidQueryParameterValue", "g1t takes Put Blob, Put Block and Put Block List."),
585 }
586 }
587 _ => azure_error(405, "UnsupportedHttpVerb", "That link is not for this."),
588 }
589}
590
591/// Finishes an upload from its parts: the blocks a list names, or the one
592/// part of a Put Blob.
593#[allow(clippy::too_many_arguments)]
594async fn complete_blob(
595 bucket: &Bucket,
596 services: &Services,
597 token: &str,
598 grant: &BlobGrant,
599 upload_id: &str,
600 ids: &[String],
601 limit: u64,
602 whole: bool,
603) -> Result<Response> {
604 let recorded: Outcome<Vec<BlobPart>> = g1t_kit::call(&services.actions, "blob_parts", &BlobArgs { blob: token.to_owned(), ..BlobArgs::default() }).await?;
605 let recorded = match recorded {
606 Outcome::Ok(recorded) => recorded,
607 Outcome::Fail(refused) => return azure_error(403, "AuthenticationFailed", &refused.message),
608 };
609 let parts = if whole {
610 recorded.into_iter().filter(|p| p.part == 1).collect()
611 } else {
612 match parts_for(ids, &recorded) {
613 Ok(parts) => parts,
614 Err(problem) => return azure_error(400, "InvalidBlockList", &problem),
615 }
616 };
617 let size = match finish(bucket, &grant.object, upload_id, &parts).await {
618 Ok(size) => size,
619 Err(problem) => return azure_error(400, "InvalidBlockList", &format!("The upload could not be completed: {problem}")),
620 };
621 if size > limit {
622 bucket.delete(&grant.object).await?;
623 return azure_error(413, "RequestBodyTooLarge", &format!("It is {} MB, more than g1t keeps ({} MB).", size / 1_048_576, limit / 1_048_576));
624 }
625 let _: Outcome<bool> = g1t_kit::call(&services.actions, "blob_done", &BlobArgs { blob: token.to_owned(), size, ..BlobArgs::default() }).await?;
626 azure(201)
627}
628
629#[cfg(test)]
630mod tests {
631 use super::*;
632
633 /// What Azure's SDK sends as block ids: base64 of a 36-character uuid
634 /// prefix and the index padded to 48 characters in all.
635 fn azure_block_id(index: u32) -> String {
636 let prefix = "4a2f0d2e-8a44-4f1b-9d55-6f1a2b3c4d5e";
637 let padded = format!("{prefix}{index:0>width$}", width = 48 - prefix.len());
638 STANDARD.encode(padded)
639 }
640
641 #[test]
642 fn block_ids_give_their_index() {
643 assert_eq!(block_index(&azure_block_id(0)), Some(0));
644 assert_eq!(block_index(&azure_block_id(17)), Some(17));
645 assert_eq!(block_index(&STANDARD.encode("no-digits")), None);
646 assert_eq!(block_index("not base64!"), None);
647 }
648
649 #[test]
650 fn a_block_list_is_read_in_order_and_matched_to_parts() {
651 // As the SDK's commitBlockList sends it.
652 let xml = format!(
653 "<?xml version=\"1.0\" encoding=\"UTF-8\" standalone=\"yes\"?><BlockList><Latest>{}</Latest><Latest>{}</Latest></BlockList>",
654 azure_block_id(0),
655 azure_block_id(1)
656 );
657 let ids = block_list(&xml);
658 assert_eq!(ids, [azure_block_id(0), azure_block_id(1)]);
659 // Sent out of order, as they are concurrently.
660 let recorded = vec![
661 BlobPart { part: 2, etag: "b".into(), size: 3 },
662 BlobPart { part: 1, etag: "a".into(), size: 8 },
663 ];
664 let parts = parts_for(&ids, &recorded).unwrap();
665 assert_eq!(parts.iter().map(|p| p.etag.as_str()).collect::<Vec<_>>(), ["a", "b"]);
666 // A block never sent, or listed out of order.
667 assert!(parts_for(&[azure_block_id(0), azure_block_id(2)], &recorded).is_err());
668 assert!(parts_for(&[azure_block_id(1), azure_block_id(0)], &recorded).is_err());
669 assert!(block_list("<BlockList></BlockList>").is_empty());
670 }
671
672 #[test]
673 fn older_protocol_chunks_are_parts() {
674 let mb32 = CACHE_PART_BYTES;
675 assert_eq!(chunk_part(&format!("bytes 0-{}/*", mb32 - 1)), Some((1, mb32)));
676 assert_eq!(chunk_part(&format!("bytes {}-{}/*", mb32 * 2, mb32 * 2 + 99)), Some((3, 100)));
677 // Not on a chunk's boundary, too long, or not a range.
678 assert_eq!(chunk_part("bytes 5-10/*"), None);
679 assert_eq!(chunk_part(&format!("bytes 0-{}/*", mb32)), None);
680 assert_eq!(chunk_part("0-10"), None);
681 }
682
683 /// The toolkit's requests, as `@actions/cache` 4 and `@actions/artifact`
684 /// 2 send them (protobuf-ts, proto field names, no defaults).
685 #[test]
686 fn twirp_requests_read_in_the_toolkits_spelling() {
687 let create_cache = json!({ "key": "node-cache-Linux-x64-npm-abc", "version": "a7f2c1e0" });
688 assert_eq!(text(&create_cache, "key"), "node-cache-Linux-x64-npm-abc");
689 let lookup = json!({ "key": "k", "restore_keys": ["k-", "x-"], "version": "v" });
690 assert_eq!(field(&lookup, "restore_keys").as_array().unwrap().len(), 2);
691 let finalize = json!({ "key": "k", "size_bytes": "1048576", "version": "v" });
692 assert_eq!(number(&finalize, "size_bytes"), Some(1_048_576));
693 let create_artifact = json!({
694 "workflow_run_backend_id": "run_1",
695 "workflow_job_run_backend_id": "job_1",
696 "name": "dist",
697 "expires_at": "2026-10-13T00:00:00Z",
698 "version": 4
699 });
700 assert_eq!(text(&create_artifact, "workflow_job_run_backend_id"), "job_1");
701 let list = json!({ "workflow_run_backend_id": "run_1", "workflow_job_run_backend_id": "job_1", "id_filter": "42", "name_filter": "dist" });
702 assert_eq!(number(&list, "id_filter"), Some(42));
703 assert_eq!(text(&list, "name_filter"), "dist");
704 // A client writing JSON names instead reads the same.
705 let camel = json!({ "workflowRunBackendId": "run_1", "sizeBytes": 3 });
706 assert_eq!(text(&camel, "workflow_run_backend_id"), "run_1");
707 assert_eq!(number(&camel, "size_bytes"), Some(3));
708 }
709
710 #[test]
711 fn the_runtime_token_names_its_run_and_job() {
712 let payload = URL_SAFE_NO_PAD.encode(json!({ "job": "job_1", "run": "run_1" }).to_string());
713 let token = format!("h.{payload}.s");
714 assert_eq!(runtime_job(&token).as_deref(), Some("job_1"));
715 assert_eq!(backend_ids(&token), ("run_1".to_owned(), "job_1".to_owned()));
716 assert_eq!(runtime_job("deadbeef"), None);
717 }
718
719 #[test]
720 fn a_job_is_told_where_the_toolkit_s_services_are() {
721 let vars = runtime_variables("https://api.g1t.sh", "tok", false);
722 // Twirp paths are resolved against the root of ACTIONS_RESULTS_URL.
723 assert_eq!(vars["ACTIONS_RESULTS_URL"], "https://api.g1t.sh/");
724 // The older protocol appends `_apis/artifactcache/…`.
725 assert_eq!(vars["ACTIONS_CACHE_URL"], "https://api.g1t.sh/actions/toolkit/");
726 assert_eq!(vars["ACTIONS_CACHE_SERVICE_V2"], "True");
727 assert!(vars.get("ACTIONS_ID_TOKEN_REQUEST_URL").is_none());
728 let vars = runtime_variables("https://api.g1t.sh", "tok", true);
729 // core.getIDToken appends `&audience=…`.
730 assert_eq!(vars["ACTIONS_ID_TOKEN_REQUEST_URL"], "https://api.g1t.sh/actions/oidc/token?api-version=2.0");
731 assert_eq!(vars["ACTIONS_ID_TOKEN_REQUEST_TOKEN"], "tok");
732 }
733
734 #[test]
735 fn artifacts_are_listed_as_the_toolkit_reads_them() {
736 let artifact = Artifact {
737 id: 7,
738 name: "dist".into(),
739 size: 10,
740 digest: None,
741 format: "zip".into(),
742 run_id: "run_1".into(),
743 job_id: "job_1".into(),
744 repo_id: "repo_1".into(),
745 expired: false,
746 created_at: "2026-10-08T12:00:00.000Z".into(),
747 updated_at: "2026-10-08T12:00:00.000Z".into(),
748 expires_at: "2026-10-22T12:00:00.000Z".into(),
749 head_branch: None,
750 head_sha: None,
751 };
752 let shown = listed(&artifact);
753 assert_eq!(shown["database_id"], "7");
754 assert_eq!(shown["size"], "10");
755 assert_eq!(shown["workflow_job_run_backend_id"], "job_1");
756 }
757}

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