Skip to content
1,225 linesCodeBlameRaw
1//! 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; the cache's methods also in protobuf (`application/protobuf`),
14//! which other clients of the protocol send (sccache, through OpenDAL).
15//! - The cache's older protocol, at `{ACTIONS_CACHE_URL}_apis/artifactcache/…`,
16//! which the toolkit's client uses whenever the server it runs against is
17//! not github.com: on g1t, that is the one it uses, and sccache's too.
18//! Its entries are sent in 32 MB chunks, or in one chunk of any size.
19//! - Blobs, at `/actions/toolkit/blobs/{token}`: the signed links those
20//! hand out. Downloads are a GET, of the whole blob or of one byte range
21//! (`Range`), as the toolkit's client fetches large entries in segments.
22//! Uploads speak the part of Azure Blob Storage's protocol the toolkit's
23//! client uses (Put Blob, Put Block, Put Block List), mapped onto an R2
24//! multipart upload: a block's id ends in its index, which is its part's
25//! number. An upload link carries a query, as an Azure SAS link does,
26//! which clients that sign their requests with it need.
27//!
28//! Every call carries the job's runtime token; the actions service checks
29//! it and keeps the entries (cache.rs, artifacts.rs, runtime.rs there).
30
31use base64::Engine;
32use base64::engine::general_purpose::{STANDARD, URL_SAFE_NO_PAD};
33use g1t_contracts::actions::{
34 ARTIFACT_MAX_BYTES, Artifact, ArtifactBlob, ArtifactCommitArgs, ArtifactReservation, ArtifactReserveArgs, BlobArgs, BlobGrant, BlobPart,
35 BlobSignArgs, CACHE_MAX_ENTRY_BYTES, CACHE_PART_BYTES, CacheCommitArgs, CacheCommitted, CacheHit, CacheLookupArgs, CacheReservation,
36 CacheReserveArgs, CacheUploadArgs, JobArtifactsArgs,
37};
38use g1t_contracts::{Failure, FailureCode, Outcome};
39use serde_json::{Map, Value, json};
40use worker::{Bucket, Env, Request, Response, Result, UploadedPart};
41
42use crate::artifacts::blob_url;
43use crate::operations::Services;
44
45/// The Twirp services, by name.
46pub const CACHE_SERVICE: &str = "github.actions.results.api.v1.CacheService";
47pub const ARTIFACT_SERVICE: &str = "github.actions.results.api.v1.ArtifactService";
48/// Where the cache's older protocol is: `ACTIONS_CACHE_URL`.
49pub const CACHE_PATH: &str = "/actions/toolkit/";
50/// The largest single block or blob a request may carry.
51const MAX_BLOCK_BYTES: u64 = 256 * 1024 * 1024;
52
53/// The bearer token of a request.
54pub fn bearer(request: &Request) -> String {
55 request
56 .headers()
57 .get("authorization")
58 .ok()
59 .flatten()
60 .and_then(|h| h.split_once(' ').map(|(_, t)| t.trim().to_owned()))
61 .unwrap_or_default()
62}
63
64fn claims(token: &str) -> Option<Value> {
65 serde_json::from_slice(&URL_SAFE_NO_PAD.decode(token.split('.').nth(1)?).ok()?).ok()
66}
67
68/// The job a runtime token names, unchecked: the actions service checks it.
69pub fn runtime_job(token: &str) -> Option<String> {
70 claims(token)?["job"].as_str().map(str::to_owned)
71}
72
73/// The variables a job gets for the toolkit: its runtime token and where
74/// the services are, and where to ask for an OIDC token when it may.
75pub fn runtime_variables(api: &str, token: &str, id_token: bool) -> Map<String, Value> {
76 let mut vars = Map::new();
77 let mut set = |k: &str, v: String| {
78 vars.insert(k.to_owned(), Value::String(v));
79 };
80 set("ACTIONS_RUNTIME_TOKEN", token.to_owned());
81 set("ACTIONS_RESULTS_URL", format!("{api}/"));
82 set("ACTIONS_CACHE_URL", format!("{api}{CACHE_PATH}"));
83 set("ACTIONS_CACHE_SERVICE_V2", "True".to_owned());
84 if id_token {
85 set("ACTIONS_ID_TOKEN_REQUEST_URL", format!("{}/token?api-version=2.0", crate::oidc::issuer(api)));
86 set("ACTIONS_ID_TOKEN_REQUEST_TOKEN", token.to_owned());
87 }
88 vars
89}
90
91// ── Twirp ───────────────────────────────────────────────────────────────────
92
93/// Twirp's binary encoding.
94const PROTOBUF: &str = "application/protobuf";
95
96/// Whether a request's `Content-Type` is Twirp's protobuf encoding.
97fn is_protobuf(content_type: &str) -> bool {
98 let kind = content_type.split(';').next().unwrap_or_default().trim().to_ascii_lowercase();
99 kind == PROTOBUF || kind == "application/x-protobuf"
100}
101
102/// The cache service's messages in protobuf, read into and written from the
103/// JSON the handlers use (`results/api/v1/cache.proto`, field numbers as
104/// there). Only what the cache's three methods carry: strings, a repeated
105/// string, an int64 and a bool. `metadata` (field 1 of each request) is
106/// skipped: the runtime token says whose cache it is.
107mod proto {
108 use serde_json::{Map, Value};
109
110 #[derive(Clone, Copy)]
111 enum Kind {
112 Text,
113 Texts,
114 Int,
115 Bool,
116 }
117
118 /// A message's fields: number, JSON name, kind.
119 type Fields = &'static [(u64, &'static str, Kind)];
120
121 fn request_fields(method: &str) -> Option<Fields> {
122 Some(match method {
123 "CreateCacheEntry" => &[(2, "key", Kind::Text), (3, "version", Kind::Text)],
124 "FinalizeCacheEntryUpload" => &[(2, "key", Kind::Text), (3, "size_bytes", Kind::Int), (4, "version", Kind::Text)],
125 "GetCacheEntryDownloadURL" => &[(2, "key", Kind::Text), (3, "restore_keys", Kind::Texts), (4, "version", Kind::Text)],
126 _ => return None,
127 })
128 }
129
130 fn response_fields(method: &str) -> Fields {
131 match method {
132 "CreateCacheEntry" => &[(1, "ok", Kind::Bool), (2, "signed_upload_url", Kind::Text), (3, "message", Kind::Text)],
133 "FinalizeCacheEntryUpload" => &[(1, "ok", Kind::Bool), (2, "entry_id", Kind::Int), (3, "message", Kind::Text)],
134 "GetCacheEntryDownloadURL" => &[(1, "ok", Kind::Bool), (2, "signed_download_url", Kind::Text), (3, "matched_key", Kind::Text)],
135 _ => &[],
136 }
137 }
138
139 fn varint(bytes: &[u8], at: &mut usize) -> Option<u64> {
140 let mut value = 0u64;
141 for shift in (0..64).step_by(7) {
142 let byte = *bytes.get(*at)?;
143 *at += 1;
144 value |= u64::from(byte & 0x7f) << shift;
145 if byte & 0x80 == 0 {
146 return Some(value);
147 }
148 }
149 None
150 }
151
152 fn put_varint(out: &mut Vec<u8>, mut value: u64) {
153 while value >= 0x80 {
154 out.push((value as u8 & 0x7f) | 0x80);
155 value >>= 7;
156 }
157 out.push(value as u8);
158 }
159
160 /// A request of `method` as JSON, or None when it is not one.
161 pub fn request(method: &str, bytes: &[u8]) -> Option<Value> {
162 let fields = request_fields(method)?;
163 let mut out = Map::new();
164 let mut at = 0;
165 while at < bytes.len() {
166 let tag = varint(bytes, &mut at)?;
167 let (number, wire) = (tag >> 3, tag & 7);
168 let known = fields.iter().find(|(n, _, _)| *n == number);
169 match wire {
170 0 => {
171 let value = varint(bytes, &mut at)?;
172 if let Some((_, name, Kind::Int)) = known {
173 // An int64 is sent as its two's complement.
174 out.insert((*name).to_owned(), Value::String((value as i64).to_string()));
175 }
176 }
177 2 => {
178 let length = usize::try_from(varint(bytes, &mut at)?).ok()?;
179 let end = at.checked_add(length).filter(|end| *end <= bytes.len())?;
180 let raw = &bytes[at..end];
181 at = end;
182 match known {
183 Some((_, name, Kind::Text)) => {
184 out.insert((*name).to_owned(), Value::String(String::from_utf8(raw.to_vec()).ok()?));
185 }
186 Some((_, name, Kind::Texts)) => {
187 let text = Value::String(String::from_utf8(raw.to_vec()).ok()?);
188 match out.entry((*name).to_owned()).or_insert_with(|| Value::Array(Vec::new())) {
189 Value::Array(list) => list.push(text),
190 _ => return None,
191 }
192 }
193 _ => {}
194 }
195 }
196 1 => at = at.checked_add(8).filter(|end| *end <= bytes.len())?,
197 5 => at = at.checked_add(4).filter(|end| *end <= bytes.len())?,
198 _ => return None,
199 }
200 }
201 Some(Value::Object(out))
202 }
203
204 /// A response of `method` from its JSON. Defaults are left out, as
205 /// proto3 does.
206 pub fn response(method: &str, value: &Value) -> Vec<u8> {
207 let mut out = Vec::new();
208 for (number, name, kind) in response_fields(method) {
209 let field = &value[*name];
210 match kind {
211 Kind::Bool if field.as_bool() == Some(true) => {
212 put_varint(&mut out, number << 3);
213 put_varint(&mut out, 1);
214 }
215 Kind::Int => {
216 let n = field.as_i64().or_else(|| field.as_str().and_then(|s| s.parse().ok())).unwrap_or(0);
217 if n != 0 {
218 put_varint(&mut out, number << 3);
219 put_varint(&mut out, n as u64);
220 }
221 }
222 Kind::Text => {
223 let text = field.as_str().unwrap_or_default();
224 if !text.is_empty() {
225 put_varint(&mut out, (number << 3) | 2);
226 put_varint(&mut out, text.len() as u64);
227 out.extend_from_slice(text.as_bytes());
228 }
229 }
230 _ => {}
231 }
232 }
233 out
234 }
235}
236
237/// A Twirp error: its code and message, at the status Twirp gives it.
238fn twirp_error(code: &str, message: &str) -> Result<Response> {
239 let status = match code {
240 "unauthenticated" => 401,
241 "permission_denied" => 403,
242 "not_found" => 404,
243 "already_exists" => 409,
244 "invalid_argument" | "malformed" => 400,
245 "bad_route" => 404,
246 "failed_precondition" => 412,
247 "resource_exhausted" => 429,
248 _ => 500,
249 };
250 Ok(Response::from_json(&json!({ "code": code, "msg": message }))?.with_status(status))
251}
252
253fn twirp_failure(failure: &Failure) -> Result<Response> {
254 let code = match failure.code {
255 FailureCode::Unauthenticated => "unauthenticated",
256 FailureCode::Forbidden => "permission_denied",
257 FailureCode::NotFound => "not_found",
258 FailureCode::Conflict => "already_exists",
259 FailureCode::Invalid => "invalid_argument",
260 _ => "failed_precondition",
261 };
262 twirp_error(code, &failure.message)
263}
264
265/// A field in the toolkit's spelling (`snake_case`), or its JSON name.
266fn field<'a>(body: &'a Value, name: &str) -> &'a Value {
267 if !body[name].is_null() {
268 return &body[name];
269 }
270 let camel: String = name.split('_').enumerate().map(|(i, p)| if i == 0 { p.to_owned() } else { p[..1].to_uppercase() + &p[1..] }).collect();
271 &body[camel]
272}
273
274fn text(body: &Value, name: &str) -> String {
275 match field(body, name) {
276 Value::String(s) => s.clone(),
277 Value::Number(n) => n.to_string(),
278 // A wrapper written as an object, `{ "value": … }`.
279 Value::Object(o) => o.get("value").map(|v| v.as_str().map_or_else(|| v.to_string(), str::to_owned)).unwrap_or_default(),
280 _ => String::new(),
281 }
282}
283
284fn number(body: &Value, name: &str) -> Option<u64> {
285 text(body, name).trim().parse().ok()
286}
287
288/// The run and job a runtime token names, which a request's backend ids
289/// must match.
290fn backend_ids(token: &str) -> (String, String) {
291 let c = claims(token).unwrap_or_default();
292 (c["run"].as_str().unwrap_or_default().to_owned(), c["job"].as_str().unwrap_or_default().to_owned())
293}
294
295/// What the toolkit's artifact client lists.
296fn listed(artifact: &Artifact) -> Value {
297 json!({
298 "workflow_run_backend_id": artifact.run_id,
299 "workflow_job_run_backend_id": artifact.job_id,
300 "database_id": artifact.id.to_string(),
301 "name": artifact.name,
302 "size": artifact.size.to_string(),
303 "created_at": artifact.created_at,
304 "digest": artifact.digest,
305 })
306}
307
308/// The Azure Storage version g1t's blob links answer as.
309const AZURE_VERSION: &str = "2024-11-04";
310
311/// An upload link: the blob's, with a query as an Azure SAS link has one.
312/// A client that treats it as a container, a blob and a SAS token (OpenDAL,
313/// which sccache uses) refuses a link without one; the token in the path
314/// is what g1t checks.
315pub fn upload_url(api: &str, blob: &str) -> String {
316 format!("{}?sv={AZURE_VERSION}", blob_url(api, blob))
317}
318
319/// Starts an R2 upload for an entry the service reserved, and the signed
320/// link the toolkit sends it to.
321async fn start_upload(bucket: &Bucket, services: &Services, job: &str, token: &str, kind: &str, id: &str, object: &str) -> Result<Outcome<String>> {
322 let upload = bucket.create_multipart_upload(object).execute().await?;
323 let upload = upload.upload_id().await;
324 let signed: Outcome<String> = g1t_kit::call(
325 &services.actions,
326 "blob_sign",
327 &BlobSignArgs { job: job.to_owned(), token: token.to_owned(), kind: kind.to_owned(), id: id.to_owned(), upload },
328 )
329 .await?;
330 Ok(match signed {
331 Outcome::Ok(blob) => Outcome::Ok(upload_url(&services.addresses.api, &blob)),
332 Outcome::Fail(refused) => Outcome::Fail(refused),
333 })
334}
335
336/// `POST /twirp/{service}/{method}`. A failure inside is logged and
337/// answered as Twirp's `internal`, with its cause.
338pub async fn twirp(request: Request, env: &Env, services: &Services, service: &str, method: &str) -> Result<Response> {
339 match twirp_inner(request, env, services, service, method).await {
340 Ok(response) => Ok(response),
341 Err(error) => twirp_error("internal", &failed(&format!("POST /twirp/{service}/{method}"), &error)),
342 }
343}
344
345async fn twirp_inner(mut request: Request, env: &Env, services: &Services, service: &str, method: &str) -> Result<Response> {
346 let token = bearer(&request);
347 let Some(job) = runtime_job(&token) else {
348 return twirp_error("unauthenticated", "Send the job's ACTIONS_RUNTIME_TOKEN as a bearer token.");
349 };
350 // Twirp clients send JSON or protobuf, and are answered in kind.
351 let binary = is_protobuf(&request.headers().get("content-type")?.unwrap_or_default());
352 let body: Value = if binary {
353 match proto::request(method, &request.bytes().await.unwrap_or_default()) {
354 Some(body) => body,
355 None => return twirp_error("malformed", "That is not a protobuf message this method takes."),
356 }
357 } else {
358 request.json().await.unwrap_or(Value::Null)
359 };
360 let answer = |value: Value| -> Result<Response> {
361 if binary {
362 let mut response = Response::from_bytes(proto::response(method, &value))?;
363 response.headers_mut().set("content-type", PROTOBUF)?;
364 Ok(response)
365 } else {
366 Response::from_json(&value)
367 }
368 };
369 let bucket = env.bucket("ACTIONS_CACHE")?;
370 let actions = &services.actions;
371 let (run, own_job) = backend_ids(&token);
372 // An artifact call names its run, and for an upload its job: the
373 // token's.
374 if service == ARTIFACT_SERVICE {
375 let asked_run = text(&body, "workflow_run_backend_id");
376 if !asked_run.is_empty() && asked_run != run {
377 return twirp_error("permission_denied", "The runtime token is for another run.");
378 }
379 let asked_job = text(&body, "workflow_job_run_backend_id");
380 if matches!(method, "CreateArtifact" | "FinalizeArtifact") && !asked_job.is_empty() && asked_job != own_job {
381 return twirp_error("permission_denied", "The runtime token is for another job.");
382 }
383 }
384 match (service, method) {
385 (CACHE_SERVICE, "GetCacheEntryDownloadURL") => {
386 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();
387 let args = CacheLookupArgs { job, token, key: text(&body, "key"), restore, version: Some(text(&body, "version")) };
388 let found: Outcome<Option<CacheHit>> = g1t_kit::call(actions, "cache_lookup", &args).await?;
389 match found {
390 Outcome::Ok(Some(CacheHit { key, blob: Some(blob), .. })) => {
391 answer(json!({ "ok": true, "signed_download_url": blob_url(&services.addresses.api, &blob), "matched_key": key }))
392 }
393 Outcome::Ok(_) => answer(json!({ "ok": false, "signed_download_url": "", "matched_key": "" })),
394 Outcome::Fail(refused) => twirp_failure(&refused),
395 }
396 }
397 (CACHE_SERVICE, "CreateCacheEntry") => {
398 let args = CacheReserveArgs { job: job.clone(), token: token.clone(), key: text(&body, "key"), size: 0, version: Some(text(&body, "version")) };
399 let reserved: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_reserve", &args).await?;
400 let reserved = match reserved {
401 Outcome::Ok(reserved) => reserved,
402 // A key already saved, or being saved by another job, to a
403 // protobuf client (OpenDAL's) is Twirp's `already_exists`
404 // (409), which it takes as "someone else has it": sccache
405 // then still writes. An `ok: false` would read as a broken
406 // cache, and sccache would only read from it.
407 Outcome::Fail(refused) if binary && refused.code == FailureCode::Conflict => return twirp_failure(&refused),
408 // The toolkit's client logs this ("another job may be
409 // creating this cache") and goes on.
410 Outcome::Fail(refused) => return answer(json!({ "ok": false, "signed_upload_url": "", "message": refused.message })),
411 };
412 match start_upload(&bucket, services, &job, &token, "cache", &reserved.id, &reserved.object).await? {
413 Outcome::Ok(url) => answer(json!({ "ok": true, "signed_upload_url": url })),
414 Outcome::Fail(refused) => answer(json!({ "ok": false, "signed_upload_url": "", "message": refused.message })),
415 }
416 }
417 (CACHE_SERVICE, "FinalizeCacheEntryUpload") => {
418 let args = CacheUploadArgs { job: job.clone(), token: token.clone(), number: None, key: Some(text(&body, "key")), version: Some(text(&body, "version")) };
419 let pending: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_upload", &args).await?;
420 let pending = match pending {
421 Outcome::Ok(pending) => pending,
422 Outcome::Fail(refused) => return answer(json!({ "ok": false, "entry_id": "0", "message": refused.message })),
423 };
424 let Some(object) = bucket.head(&pending.object).await? else {
425 return answer(json!({ "ok": false, "entry_id": "0", "message": "Nothing was uploaded for that entry." }));
426 };
427 match commit_cache(&bucket, services, &job, &token, &pending.id, object.size()).await? {
428 Outcome::Ok(()) => answer(json!({ "ok": true, "entry_id": pending.number.to_string() })),
429 Outcome::Fail(refused) => answer(json!({ "ok": false, "entry_id": "0", "message": refused.message })),
430 }
431 }
432 (ARTIFACT_SERVICE, "CreateArtifact") => {
433 let expires_at = Some(text(&body, "expires_at")).filter(|e| !e.is_empty());
434 let args = ArtifactReserveArgs { job: job.clone(), token: token.clone(), name: text(&body, "name"), expires_at, ..ArtifactReserveArgs::default() };
435 let reserved: Outcome<ArtifactReservation> = g1t_kit::call(actions, "artifact_reserve", &args).await?;
436 let reserved = match reserved {
437 Outcome::Ok(reserved) => reserved,
438 Outcome::Fail(refused) => return twirp_failure(&refused),
439 };
440 match start_upload(&bucket, services, &job, &token, "artifact", &reserved.id.to_string(), &reserved.object).await? {
441 Outcome::Ok(url) => Response::from_json(&json!({ "ok": true, "signed_upload_url": url })),
442 Outcome::Fail(refused) => twirp_failure(&refused),
443 }
444 }
445 (ARTIFACT_SERVICE, "FinalizeArtifact") => {
446 let digest = Some(text(&body, "hash")).filter(|h| !h.is_empty());
447 // The size it was measured at as it was stored, not the one it
448 // says.
449 let args = ArtifactCommitArgs { job, token, id: None, name: Some(text(&body, "name")), size: 0, digest };
450 let done: Outcome<Artifact> = g1t_kit::call(actions, "artifact_commit", &args).await?;
451 match done {
452 Outcome::Ok(artifact) => Response::from_json(&json!({ "ok": true, "artifact_id": artifact.id.to_string() })),
453 Outcome::Fail(refused) => twirp_failure(&refused),
454 }
455 }
456 (ARTIFACT_SERVICE, "ListArtifacts") => {
457 let name = Some(text(&body, "name_filter")).filter(|n| !n.is_empty());
458 let id = number(&body, "id_filter");
459 let args = JobArtifactsArgs { job, token, run_id: None, name, id };
460 let found: Outcome<Vec<Artifact>> = g1t_kit::call(actions, "job_artifacts", &args).await?;
461 match found {
462 Outcome::Ok(found) => Response::from_json(&json!({ "artifacts": found.iter().map(listed).collect::<Vec<_>>() })),
463 Outcome::Fail(refused) => twirp_failure(&refused),
464 }
465 }
466 (ARTIFACT_SERVICE, "GetSignedArtifactURL") => {
467 let args = JobArtifactsArgs { job, token, run_id: None, name: Some(text(&body, "name")), id: None };
468 let found: Outcome<ArtifactBlob> = g1t_kit::call(actions, "job_artifact", &args).await?;
469 match found {
470 Outcome::Ok(found) if !found.blob.is_empty() => Response::from_json(&json!({ "signed_url": blob_url(&services.addresses.api, &found.blob) })),
471 Outcome::Ok(_) => twirp_error("failed_precondition", "Download links are not set up on this installation."),
472 Outcome::Fail(refused) => twirp_failure(&refused),
473 }
474 }
475 (ARTIFACT_SERVICE, "DeleteArtifact") => {
476 let args = JobArtifactsArgs { job, token, run_id: None, name: Some(text(&body, "name")), id: None };
477 let done: Outcome<Artifact> = g1t_kit::call(actions, "job_delete_artifact", &args).await?;
478 match done {
479 Outcome::Ok(artifact) => Response::from_json(&json!({ "ok": true, "artifact_id": artifact.id.to_string() })),
480 Outcome::Fail(refused) => twirp_failure(&refused),
481 }
482 }
483 _ => twirp_error("bad_route", &format!("No method {method} on {service}.")),
484 }
485}
486
487/// Marks an uploaded cache entry ready, and deletes what that evicted.
488async fn commit_cache(bucket: &Bucket, services: &Services, job: &str, token: &str, id: &str, size: u64) -> Result<Outcome<()>> {
489 let committed: Outcome<CacheCommitted> = g1t_kit::call(
490 &services.actions,
491 "cache_commit",
492 &CacheCommitArgs { job: job.to_owned(), token: token.to_owned(), id: id.to_owned(), size },
493 )
494 .await?;
495 Ok(match committed {
496 Outcome::Ok(committed) => {
497 if !committed.evicted.is_empty() {
498 bucket.delete_multiple(committed.evicted.iter().map(String::as_str).collect()).await?;
499 }
500 Outcome::Ok(())
501 }
502 Outcome::Fail(refused) => Outcome::Fail(refused),
503 })
504}
505
506// ── The cache's older protocol ──────────────────────────────────────────────
507
508fn plain_error(status: u16, message: &str) -> Result<Response> {
509 Ok(Response::from_json(&json!({ "message": message, "error": { "message": message } }))?.with_status(status))
510}
511
512fn query(request: &Request, name: &str) -> Option<String> {
513 query_in(&request.url().ok()?, name)
514}
515
516fn query_in(url: &worker::Url, name: &str) -> Option<String> {
517 url.query_pairs().find(|(k, _)| k == name).map(|(_, v)| v.into_owned())
518}
519
520/// The part a chunk of the older protocol is, from its `Content-Range`:
521/// chunks are `CACHE_PART_BYTES` apart, as the toolkit sends them. A first
522/// chunk may be larger, up to `MAX_BLOCK_BYTES`: a client that sends an
523/// entry in one request (sccache does) sends a single chunk from 0.
524pub fn chunk_part(range: &str) -> Option<(u16, u64)> {
525 let range = range.trim().strip_prefix("bytes ")?;
526 let (span, _) = range.split_once('/')?;
527 let (start, end) = span.split_once('-')?;
528 let (start, end): (u64, u64) = (start.trim().parse().ok()?, end.trim().parse().ok()?);
529 if end < start {
530 return None;
531 }
532 let length = end - start + 1;
533 if start == 0 && length <= MAX_BLOCK_BYTES {
534 return Some((1, length));
535 }
536 if start % CACHE_PART_BYTES != 0 || length > CACHE_PART_BYTES {
537 return None;
538 }
539 Some(((start / CACHE_PART_BYTES + 1) as u16, length))
540}
541
542/// The bytes a download's `Range` header asks for, out of `size`: first and
543/// last, inclusive. None to send the whole blob (no header, or one this
544/// does not read, such as several ranges); `Some(None)` when the range is
545/// past the end (416).
546pub fn byte_range(header: &str, size: u64) -> Option<Option<(u64, u64)>> {
547 let spec = header.trim().strip_prefix("bytes=")?.trim();
548 if spec.contains(',') {
549 return None;
550 }
551 let (first, last) = spec.split_once('-')?;
552 let (first, last) = (first.trim(), last.trim());
553 let range = if first.is_empty() {
554 // The last `n` bytes.
555 let n: u64 = last.parse().ok()?;
556 if n == 0 || size == 0 {
557 return Some(None);
558 }
559 (size.saturating_sub(n), size - 1)
560 } else {
561 let first: u64 = first.parse().ok()?;
562 let last: u64 = if last.is_empty() { u64::MAX } else { last.parse().ok()? };
563 if last < first {
564 return None;
565 }
566 if first >= size {
567 return Some(None);
568 }
569 (first, last.min(size - 1))
570 };
571 Some(Some(range))
572}
573
574/// What a lookup of the older protocol answers: 200 with the entry, 204
575/// for a miss (which the toolkit's client and sccache read as "not
576/// cached"), or the refusal's status.
577pub fn lookup_answer(found: Outcome<Option<CacheHit>>, version: &str, api: &str) -> (u16, Option<Value>) {
578 match found {
579 Outcome::Ok(Some(CacheHit { key, blob: Some(blob), created_at, .. })) => (
580 200,
581 Some(json!({
582 "cacheKey": key,
583 "cacheVersion": version,
584 "scope": "",
585 "creationTime": created_at,
586 "archiveLocation": blob_url(api, &blob),
587 })),
588 ),
589 // No entry, or one without a download link (no ACTIONS_KEY): a miss.
590 Outcome::Ok(_) => (204, None),
591 Outcome::Fail(refused) => (refused.code.http_status(), Some(json!({ "message": refused.message, "error": { "message": refused.message } }))),
592 }
593}
594
595/// Logs a toolkit request that failed inside g1t, and says what to tell
596/// its client: the cause, so a job's log shows more than a bare 500.
597fn failed(route: &str, error: &worker::Error) -> String {
598 worker::console_error!("toolkit: {route} failed: {error}");
599 format!("g1t could not answer this: {error}")
600}
601
602/// `{ACTIONS_CACHE_URL}_apis/artifactcache/…`. `rest` is the path after it.
603/// A failure inside is logged and answered as a 500 with its cause.
604pub async fn cache_v1(request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
605 match cache_v1_inner(request, env, services, method, rest).await {
606 Ok(response) => Ok(response),
607 Err(error) => plain_error(500, &failed(&format!("{method} {CACHE_PATH}_apis/artifactcache/{rest}"), &error)),
608 }
609}
610
611async fn cache_v1_inner(mut request: Request, env: &Env, services: &Services, method: &str, rest: &str) -> Result<Response> {
612 let token = bearer(&request);
613 let Some(job) = runtime_job(&token) else {
614 return plain_error(401, "Send the job's ACTIONS_RUNTIME_TOKEN as a bearer token.");
615 };
616 let bucket = env.bucket("ACTIONS_CACHE")?;
617 let actions = &services.actions;
618 let parts: Vec<&str> = rest.split('/').filter(|p| !p.is_empty()).collect();
619 match (method, parts.as_slice()) {
620 ("GET", ["cache"]) => {
621 let keys: Vec<String> = query(&request, "keys").unwrap_or_default().split(',').map(|k| k.trim().to_owned()).filter(|k| !k.is_empty()).collect();
622 let Some((key, restore)) = keys.split_first() else {
623 return plain_error(400, "Give keys.");
624 };
625 let version = query(&request, "version").unwrap_or_default();
626 let args = CacheLookupArgs { job, token, key: key.clone(), restore: restore.to_vec(), version: Some(version.clone()) };
627 let found: Outcome<Option<CacheHit>> = g1t_kit::call(actions, "cache_lookup", &args).await?;
628 match lookup_answer(found, &version, &services.addresses.api) {
629 (status, Some(body)) => Ok(Response::from_json(&body)?.with_status(status)),
630 (status, None) => Ok(Response::empty()?.with_status(status)),
631 }
632 }
633 ("POST", ["caches"]) => {
634 let body: Value = request.json().await.unwrap_or(Value::Null);
635 let size = body["cacheSize"].as_u64().unwrap_or(0);
636 let key = body["key"].as_str().unwrap_or_default().to_owned();
637 let version = body["version"].as_str().unwrap_or_default().to_owned();
638 let args = CacheReserveArgs { job: job.clone(), token: token.clone(), key, size, version: Some(version) };
639 let reserved: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_reserve", &args).await?;
640 let reserved = match reserved {
641 Outcome::Ok(reserved) => reserved,
642 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
643 };
644 match start_upload(&bucket, services, &job, &token, "cache", &reserved.id, &reserved.object).await? {
645 Outcome::Ok(_) => Ok(Response::from_json(&json!({ "cacheId": reserved.number }))?.with_status(201)),
646 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
647 }
648 }
649 (_, ["caches", number]) => {
650 let Ok(number) = number.parse::<u64>() else {
651 return plain_error(404, "No such cache entry.");
652 };
653 let args = CacheUploadArgs { job: job.clone(), token: token.clone(), number: Some(number), key: None, version: None };
654 let pending: Outcome<CacheReservation> = g1t_kit::call(actions, "cache_upload", &args).await?;
655 let pending = match pending {
656 Outcome::Ok(pending) => pending,
657 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
658 };
659 let (Some(blob), Some(upload_id)) = (pending.blob.clone(), pending.upload.clone()) else {
660 return plain_error(409, "That entry's upload was not started.");
661 };
662 match method {
663 "PATCH" => {
664 let range = request.headers().get("content-range")?.unwrap_or_default();
665 let Some((part, length)) = chunk_part(&range) else {
666 return plain_error(400, &format!("Send the entry in chunks of {} MB, each with its Content-Range.", CACHE_PART_BYTES / 1_048_576));
667 };
668 let Some(body) = request.inner().body() else { return plain_error(400, "The chunk is empty.") };
669 let upload = bucket.resume_multipart_upload(&pending.object, &upload_id)?;
670 let uploaded = upload.upload_part(part, body).await?;
671 let recorded = BlobArgs { blob, part: u32::from(part), etag: uploaded.etag(), size: length };
672 let _: Outcome<bool> = g1t_kit::call(actions, "blob_part", &recorded).await?;
673 Ok(Response::empty()?.with_status(204))
674 }
675 "POST" => {
676 let parts: Outcome<Vec<BlobPart>> = g1t_kit::call(actions, "blob_parts", &BlobArgs { blob: blob.clone(), ..BlobArgs::default() }).await?;
677 let parts = match parts {
678 Outcome::Ok(parts) => parts,
679 Outcome::Fail(refused) => return plain_error(refused.code.http_status(), &refused.message),
680 };
681 let size = match finish(&bucket, &pending.object, &upload_id, &parts).await {
682 Ok(size) => size,
683 Err(problem) => return plain_error(400, &format!("The entry could not be completed: {problem}")),
684 };
685 let _: Outcome<bool> = g1t_kit::call(actions, "blob_done", &BlobArgs { blob, size, ..BlobArgs::default() }).await?;
686 match commit_cache(&bucket, services, &job, &token, &pending.id, size).await? {
687 Outcome::Ok(()) => Ok(Response::empty()?.with_status(204)),
688 Outcome::Fail(refused) => plain_error(refused.code.http_status(), &refused.message),
689 }
690 }
691 _ => plain_error(405, "PATCH a chunk, or POST to commit."),
692 }
693 }
694 _ => plain_error(404, "No such endpoint."),
695 }
696}
697
698/// Completes an R2 upload from its recorded parts, in order; the size it
699/// came to. An upload with no parts is an empty object.
700async fn finish(bucket: &Bucket, object: &str, upload_id: &str, parts: &[BlobPart]) -> std::result::Result<u64, String> {
701 let upload = bucket.resume_multipart_upload(object, upload_id).map_err(|e| e.to_string())?;
702 if parts.is_empty() {
703 let _ = upload.abort().await;
704 bucket.put(object, Vec::<u8>::new()).execute().await.map_err(|e| e.to_string())?;
705 return Ok(0);
706 }
707 let done = upload
708 .complete(parts.iter().map(|p| UploadedPart::new(p.part as u16, p.etag.clone())))
709 .await
710 .map_err(|e| e.to_string())?;
711 Ok(done.size())
712}
713
714// ── Blobs ───────────────────────────────────────────────────────────────────
715
716/// A block's index, from its id: the toolkit's client (Azure's SDK) makes
717/// block ids as base64 of a prefix and the index padded with zeros.
718pub fn block_index(id: &str) -> Option<u32> {
719 let decoded = STANDARD.decode(id.trim()).ok()?;
720 let text = String::from_utf8(decoded).ok()?;
721 let digits: String = text.chars().rev().take_while(char::is_ascii_digit).collect::<Vec<_>>().into_iter().rev().collect();
722 if digits.is_empty() {
723 return None;
724 }
725 digits.parse().ok()
726}
727
728/// The block ids of a Put Block List body, in order.
729pub fn block_list(xml: &str) -> Vec<String> {
730 let mut ids = Vec::new();
731 let mut rest = xml;
732 while let Some(open) = rest.find('<') {
733 rest = &rest[open + 1..];
734 let Some(close) = rest.find('>') else { break };
735 let tag = &rest[..close];
736 rest = &rest[close + 1..];
737 if matches!(tag, "Latest" | "Committed" | "Uncommitted") {
738 let Some(end) = rest.find("</") else { break };
739 ids.push(rest[..end].trim().to_owned());
740 rest = &rest[end..];
741 }
742 }
743 ids
744}
745
746/// The parts a block list names, as recorded: each block must have been
747/// sent, as the part its index says, and they must run from the first.
748pub fn parts_for(ids: &[String], recorded: &[BlobPart]) -> std::result::Result<Vec<BlobPart>, String> {
749 let mut out = Vec::with_capacity(ids.len());
750 for (position, id) in ids.iter().enumerate() {
751 let index = block_index(id).ok_or_else(|| format!("The block id {id} does not end in its index."))?;
752 let part = index + 1;
753 if part as usize != position + 1 {
754 return Err("The blocks must be listed in the order they were numbered.".to_owned());
755 }
756 let found = recorded.iter().find(|p| p.part == part).ok_or_else(|| format!("Block {id} was never sent."))?;
757 out.push(found.clone());
758 }
759 Ok(out)
760}
761
762fn azure(status: u16) -> Result<Response> {
763 let mut response = Response::empty()?.with_status(status);
764 let headers = response.headers_mut();
765 headers.set("x-ms-request-id", &g1t_contracts::new_id("req", g1t_kit::now_ms()))?;
766 headers.set("x-ms-version", "2024-11-04")?;
767 headers.set("x-ms-request-server-encrypted", "true")?;
768 Ok(response)
769}
770
771fn azure_error(status: u16, code: &str, message: &str) -> Result<Response> {
772 let body = format!("<?xml version=\"1.0\" encoding=\"utf-8\"?><Error><Code>{code}</Code><Message>{message}</Message></Error>");
773 let mut response = Response::ok(body)?.with_status(status);
774 response.headers_mut().set("content-type", "application/xml")?;
775 response.headers_mut().set("x-ms-error-code", code)?;
776 Ok(response)
777}
778
779/// `/actions/toolkit/blobs/{token}`: GET or HEAD a download, PUT an upload.
780/// A failure inside is logged and answered as Azure's `InternalError`.
781pub async fn blob(request: Request, env: &Env, services: &Services, method: &str, token: &str) -> Result<Response> {
782 match blob_inner(request, env, services, method, token).await {
783 Ok(response) => Ok(response),
784 // The token is a credential: the route is logged without it.
785 Err(error) => azure_error(500, "InternalError", &failed(&format!("{method} /actions/toolkit/blobs/…"), &error)),
786 }
787}
788
789async fn blob_inner(mut request: Request, env: &Env, services: &Services, method: &str, token: &str) -> Result<Response> {
790 let opened: Outcome<BlobGrant> = g1t_kit::call(&services.actions, "blob_open", &BlobArgs { blob: token.to_owned(), ..BlobArgs::default() }).await?;
791 let grant = match opened {
792 Outcome::Ok(grant) => grant,
793 Outcome::Fail(refused) => {
794 let status = if refused.code == FailureCode::NotFound { 404 } else { 403 };
795 return azure_error(status, if status == 404 { "BlobNotFound" } else { "AuthenticationFailed" }, &refused.message);
796 }
797 };
798 let bucket = env.bucket("ACTIONS_CACHE")?;
799 match (method, grant.upload.as_deref()) {
800 ("GET" | "HEAD", None) => {
801 let headers = |response: &mut Response, size: u64| -> Result<()> {
802 let headers = response.headers_mut();
803 headers.set("content-length", &size.to_string())?;
804 headers.set("content-type", grant.content_type.as_deref().unwrap_or("application/octet-stream"))?;
805 headers.set("x-ms-blob-type", "BlockBlob")?;
806 headers.set("accept-ranges", "bytes")?;
807 if let Some(name) = &grant.filename {
808 headers.set("content-disposition", &format!("attachment; filename=\"{}\"", name.replace('"', "")))?;
809 }
810 Ok(())
811 };
812 if method == "HEAD" {
813 let Some(object) = bucket.head(&grant.object).await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
814 let mut response = Response::empty()?;
815 headers(&mut response, object.size())?;
816 return Ok(response);
817 }
818 // One byte range (`Range`, or Azure's `x-ms-range`): the
819 // toolkit's client fetches a large entry in segments, side by
820 // side, and writes each where its range says.
821 let asked = match request.headers().get("x-ms-range")? {
822 Some(range) => Some(range),
823 None => request.headers().get("range")?,
824 };
825 if let Some(asked) = asked.filter(|r| !r.trim().is_empty()) {
826 let Some(object) = bucket.head(&grant.object).await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
827 let size = object.size();
828 match byte_range(&asked, size) {
829 Some(Some((first, last))) => {
830 let length = last - first + 1;
831 let Some(object) = bucket.get(&grant.object).range(worker::Range::OffsetWithLength { offset: first, length }).execute().await? else {
832 return azure_error(404, "BlobNotFound", "It is gone.");
833 };
834 let Some(body) = object.body() else { return azure_error(404, "BlobNotFound", "It is gone.") };
835 let mut response = Response::from_body(body.response_body()?)?.with_status(206);
836 headers(&mut response, length)?;
837 response.headers_mut().set("content-range", &format!("bytes {first}-{last}/{size}"))?;
838 return Ok(response);
839 }
840 Some(None) => {
841 let mut response = azure_error(416, "InvalidRange", "The range is past the end of the blob.")?;
842 response.headers_mut().set("content-range", &format!("bytes */{size}"))?;
843 return Ok(response);
844 }
845 // Not a range this reads: the whole blob.
846 None => {}
847 }
848 }
849 let Some(object) = bucket.get(&grant.object).execute().await? else { return azure_error(404, "BlobNotFound", "It is gone.") };
850 let size = object.size();
851 let Some(body) = object.body() else { return azure_error(404, "BlobNotFound", "It is gone.") };
852 let mut response = Response::from_body(body.response_body()?)?;
853 headers(&mut response, size)?;
854 Ok(response)
855 }
856 ("PUT", Some(upload_id)) => {
857 let comp = query(&request, "comp").unwrap_or_default();
858 let limit = if grant.kind == "cache" { CACHE_MAX_ENTRY_BYTES } else { ARTIFACT_MAX_BYTES };
859 match comp.as_str() {
860 // Put Block, or Put Blob: one part.
861 "block" | "" => {
862 let part = if comp == "block" {
863 match query(&request, "blockid").as_deref().and_then(block_index) {
864 Some(index) if index < 10_000 => index + 1,
865 _ => return azure_error(400, "InvalidQueryParameterValue", "A block id ends in its index, from 0."),
866 }
867 } else {
868 1
869 };
870 let length = request.headers().get("content-length")?.and_then(|l| l.parse::<u64>().ok()).unwrap_or(0);
871 if length > MAX_BLOCK_BYTES || length > limit {
872 return azure_error(413, "RequestBodyTooLarge", "That block is larger than g1t takes at once.");
873 }
874 let upload = bucket.resume_multipart_upload(&grant.object, upload_id)?;
875 let uploaded = match request.inner().body() {
876 Some(body) if length > 0 => upload.upload_part(part as u16, body).await?,
877 _ => upload.upload_part(part as u16, Vec::<u8>::new()).await?,
878 };
879 let recorded = BlobArgs { blob: token.to_owned(), part, etag: uploaded.etag(), size: length };
880 let _: Outcome<bool> = g1t_kit::call(&services.actions, "blob_part", &recorded).await?;
881 if comp == "block" {
882 return azure(201);
883 }
884 // Put Blob is the whole thing: finish it now.
885 complete_blob(&bucket, services, token, &grant, upload_id, &[], limit, true).await
886 }
887 "blocklist" => {
888 let xml = request.text().await.unwrap_or_default();
889 let ids = block_list(&xml);
890 complete_blob(&bucket, services, token, &grant, upload_id, &ids, limit, false).await
891 }
892 _ => azure_error(400, "InvalidQueryParameterValue", "g1t takes Put Blob, Put Block and Put Block List."),
893 }
894 }
895 _ => azure_error(405, "UnsupportedHttpVerb", "That link is not for this."),
896 }
897}
898
899/// Finishes an upload from its parts: the blocks a list names, or the one
900/// part of a Put Blob.
901#[allow(clippy::too_many_arguments)]
902async fn complete_blob(
903 bucket: &Bucket,
904 services: &Services,
905 token: &str,
906 grant: &BlobGrant,
907 upload_id: &str,
908 ids: &[String],
909 limit: u64,
910 whole: bool,
911) -> Result<Response> {
912 let recorded: Outcome<Vec<BlobPart>> = g1t_kit::call(&services.actions, "blob_parts", &BlobArgs { blob: token.to_owned(), ..BlobArgs::default() }).await?;
913 let recorded = match recorded {
914 Outcome::Ok(recorded) => recorded,
915 Outcome::Fail(refused) => return azure_error(403, "AuthenticationFailed", &refused.message),
916 };
917 let parts = if whole {
918 recorded.into_iter().filter(|p| p.part == 1).collect()
919 } else {
920 match parts_for(ids, &recorded) {
921 Ok(parts) => parts,
922 Err(problem) => return azure_error(400, "InvalidBlockList", &problem),
923 }
924 };
925 let size = match finish(bucket, &grant.object, upload_id, &parts).await {
926 Ok(size) => size,
927 Err(problem) => return azure_error(400, "InvalidBlockList", &format!("The upload could not be completed: {problem}")),
928 };
929 if size > limit {
930 bucket.delete(&grant.object).await?;
931 return azure_error(413, "RequestBodyTooLarge", &format!("It is {} MB, more than g1t keeps ({} MB).", size / 1_048_576, limit / 1_048_576));
932 }
933 let _: Outcome<bool> = g1t_kit::call(&services.actions, "blob_done", &BlobArgs { blob: token.to_owned(), size, ..BlobArgs::default() }).await?;
934 azure(201)
935}
936
937#[cfg(test)]
938mod tests {
939 use super::*;
940
941 /// What Azure's SDK sends as block ids: base64 of a 36-character uuid
942 /// prefix and the index padded to 48 characters in all.
943 fn azure_block_id(index: u32) -> String {
944 let prefix = "4a2f0d2e-8a44-4f1b-9d55-6f1a2b3c4d5e";
945 let padded = format!("{prefix}{index:0>width$}", width = 48 - prefix.len());
946 STANDARD.encode(padded)
947 }
948
949 #[test]
950 fn block_ids_give_their_index() {
951 assert_eq!(block_index(&azure_block_id(0)), Some(0));
952 assert_eq!(block_index(&azure_block_id(17)), Some(17));
953 assert_eq!(block_index(&STANDARD.encode("no-digits")), None);
954 assert_eq!(block_index("not base64!"), None);
955 }
956
957 #[test]
958 fn a_block_list_is_read_in_order_and_matched_to_parts() {
959 // As the SDK's commitBlockList sends it.
960 let xml = format!(
961 "<?xml version=\"1.0\" encoding=\"UTF-8\" standalone=\"yes\"?><BlockList><Latest>{}</Latest><Latest>{}</Latest></BlockList>",
962 azure_block_id(0),
963 azure_block_id(1)
964 );
965 let ids = block_list(&xml);
966 assert_eq!(ids, [azure_block_id(0), azure_block_id(1)]);
967 // Sent out of order, as they are concurrently.
968 let recorded = vec![
969 BlobPart { part: 2, etag: "b".into(), size: 3 },
970 BlobPart { part: 1, etag: "a".into(), size: 8 },
971 ];
972 let parts = parts_for(&ids, &recorded).unwrap();
973 assert_eq!(parts.iter().map(|p| p.etag.as_str()).collect::<Vec<_>>(), ["a", "b"]);
974 // A block never sent, or listed out of order.
975 assert!(parts_for(&[azure_block_id(0), azure_block_id(2)], &recorded).is_err());
976 assert!(parts_for(&[azure_block_id(1), azure_block_id(0)], &recorded).is_err());
977 assert!(block_list("<BlockList></BlockList>").is_empty());
978 }
979
980 #[test]
981 fn older_protocol_chunks_are_parts() {
982 let mb32 = CACHE_PART_BYTES;
983 assert_eq!(chunk_part(&format!("bytes 0-{}/*", mb32 - 1)), Some((1, mb32)));
984 assert_eq!(chunk_part(&format!("bytes {}-{}/*", mb32 * 2, mb32 * 2 + 99)), Some((3, 100)));
985 // A whole entry in one chunk, as sccache (OpenDAL) sends it: its
986 // check file is 13 bytes, a compiled crate can be well over 32 MB.
987 assert_eq!(chunk_part("bytes 0-12/*"), Some((1, 13)));
988 assert_eq!(chunk_part(&format!("bytes 0-{}/*", mb32)), Some((1, mb32 + 1)));
989 assert_eq!(chunk_part(&format!("bytes 0-{}/*", MAX_BLOCK_BYTES - 1)), Some((1, MAX_BLOCK_BYTES)));
990 // Not on a chunk's boundary, too long, backwards, or not a range.
991 assert_eq!(chunk_part("bytes 5-10/*"), None);
992 assert_eq!(chunk_part(&format!("bytes {mb32}-{}/*", mb32 * 2)), None);
993 assert_eq!(chunk_part(&format!("bytes 0-{}/*", MAX_BLOCK_BYTES)), None);
994 assert_eq!(chunk_part("bytes 10-5/*"), None);
995 assert_eq!(chunk_part("0-10"), None);
996 }
997
998 #[test]
999 fn downloads_read_one_byte_range() {
1000 // OpenDAL's stat: the first byte.
1001 assert_eq!(byte_range("bytes=0-0", 100), Some(Some((0, 0))));
1002 // The toolkit's segments, the last one cut at the end.
1003 assert_eq!(byte_range("bytes=0-49", 100), Some(Some((0, 49))));
1004 assert_eq!(byte_range("bytes=50-999", 100), Some(Some((50, 99))));
1005 assert_eq!(byte_range("bytes=90-", 100), Some(Some((90, 99))));
1006 assert_eq!(byte_range("bytes=-10", 100), Some(Some((90, 99))));
1007 assert_eq!(byte_range("bytes=-1000", 100), Some(Some((0, 99))));
1008 // Past the end: 416.
1009 assert_eq!(byte_range("bytes=100-200", 100), Some(None));
1010 assert_eq!(byte_range("bytes=0-0", 0), Some(None));
1011 assert_eq!(byte_range("bytes=-0", 100), Some(None));
1012 // Not read: the whole blob.
1013 assert_eq!(byte_range("bytes=0-1,5-6", 100), None);
1014 assert_eq!(byte_range("bytes=9-3", 100), None);
1015 assert_eq!(byte_range("items=0-1", 100), None);
1016 assert_eq!(byte_range("bytes=a-b", 100), None);
1017 }
1018
1019 #[test]
1020 fn upload_links_carry_a_query_as_sas_links_do() {
1021 let url = upload_url("https://api.g1t.sh", "tok.sig");
1022 assert_eq!(url, "https://api.g1t.sh/actions/toolkit/blobs/tok.sig?sv=2024-11-04");
1023 // How OpenDAL reads a signed upload link: a container, a blob in
1024 // it, and a SAS query, all of which must be there.
1025 let rest = url.strip_prefix("https://api.g1t.sh/").unwrap();
1026 let (path, query) = rest.split_once('?').unwrap();
1027 let (container, blob) = path.split_once('/').unwrap();
1028 assert_eq!((container, blob, query), ("actions", "toolkit/blobs/tok.sig", "sv=2024-11-04"));
1029 }
1030
1031 /// A protobuf length-delimited field, as prost writes it.
1032 fn pb_text(number: u8, text: &str) -> Vec<u8> {
1033 let mut out = vec![(number << 3) | 2, text.len() as u8];
1034 out.extend_from_slice(text.as_bytes());
1035 out
1036 }
1037
1038 /// The requests sccache 0.18 sends (OpenDAL 0.58's `ghac` service, with
1039 /// prost): fields in number order, defaults left out, no metadata.
1040 #[test]
1041 fn twirp_reads_sccaches_protobuf_requests() {
1042 assert!(is_protobuf("application/protobuf"));
1043 assert!(is_protobuf("Application/Protobuf; charset=utf-8"));
1044 assert!(!is_protobuf("application/json"));
1045 assert!(!is_protobuf(""));
1046
1047 let key = "sccache/f/c/b/fcb0a1d2e3";
1048 let version = "sccache-v0.18.0";
1049 let create = [pb_text(2, key), pb_text(3, version)].concat();
1050 let read = proto::request("CreateCacheEntry", &create).unwrap();
1051 assert_eq!((text(&read, "key"), text(&read, "version")), (key.to_owned(), version.to_owned()));
1052
1053 // size_bytes is field 3, a varint: 300 is 0xac 0x02.
1054 let finalize = [pb_text(2, key), vec![0x18, 0xac, 0x02], pb_text(4, version)].concat();
1055 let read = proto::request("FinalizeCacheEntryUpload", &finalize).unwrap();
1056 assert_eq!(number(&read, "size_bytes"), Some(300));
1057 assert_eq!(text(&read, "version"), version);
1058
1059 let lookup = [pb_text(2, key), pb_text(4, version)].concat();
1060 let read = proto::request("GetCacheEntryDownloadURL", &lookup).unwrap();
1061 assert_eq!(text(&read, "key"), key);
1062 assert!(field(&read, "restore_keys").is_null());
1063 // The toolkit's own lookup, with metadata (skipped) and restore keys.
1064 let metadata = vec![0x0a, 0x02, 0x08, 0x07];
1065 let with_restore = [metadata, pb_text(2, "k"), pb_text(3, "k-"), pb_text(3, "x-"), pb_text(4, "v")].concat();
1066 let read = proto::request("GetCacheEntryDownloadURL", &with_restore).unwrap();
1067 assert_eq!(field(&read, "restore_keys"), &json!(["k-", "x-"]));
1068 assert_eq!(text(&read, "version"), "v");
1069
1070 // The same three, as prost 0.14 encodes them with OpenDAL's
1071 // generated types (the `ghac` crate, 0.3.0), byte for byte.
1072 let recorded = |hex: &str| -> Vec<u8> { (0..hex.len()).step_by(2).map(|i| u8::from_str_radix(&hex[i..i + 2], 16).unwrap()).collect() };
1073 let prefix = "1218736363616368652f662f632f622f66636230613164326533";
1074 let suffix = "0f736363616368652d76302e31382e30";
1075 assert_eq!(recorded(&format!("{prefix}1a{suffix}")), create);
1076 assert_eq!(recorded(&format!("{prefix}18ac0222{suffix}")), finalize);
1077 assert_eq!(recorded(&format!("{prefix}22{suffix}")), lookup);
1078
1079 // Cut short, or not a cache method.
1080 assert!(proto::request("CreateCacheEntry", &create[..create.len() - 1]).is_none());
1081 assert!(proto::request("CreateCacheEntry", &[0x12, 0xff]).is_none());
1082 assert!(proto::request("CreateArtifact", &create).is_none());
1083 assert_eq!(proto::request("CreateCacheEntry", &[]), Some(json!({})));
1084 }
1085
1086 #[test]
1087 fn twirp_answers_in_protobuf_as_prost_reads_it() {
1088 let url = "https://api.g1t.sh/actions/toolkit/blobs/t?sv=2024-11-04";
1089 let created = proto::response("CreateCacheEntry", &json!({ "ok": true, "signed_upload_url": url }));
1090 assert_eq!(created, [vec![0x08, 0x01], pb_text(2, url)].concat());
1091 // Refused: ok false is the default, so only the message is sent.
1092 let refused = proto::response("CreateCacheEntry", &json!({ "ok": false, "signed_upload_url": "", "message": "no" }));
1093 assert_eq!(refused, pb_text(3, "no"));
1094 // entry_id is an int64, given as a string in JSON.
1095 let finalized = proto::response("FinalizeCacheEntryUpload", &json!({ "ok": true, "entry_id": "300" }));
1096 assert_eq!(finalized, vec![0x08, 0x01, 0x10, 0xac, 0x02]);
1097 let found = proto::response("GetCacheEntryDownloadURL", &json!({ "ok": true, "signed_download_url": "u", "matched_key": "k" }));
1098 assert_eq!(found, [vec![0x08, 0x01], pb_text(2, "u"), pb_text(3, "k")].concat());
1099 // A miss is an empty message: ok false.
1100 assert!(proto::response("GetCacheEntryDownloadURL", &json!({ "ok": false, "signed_download_url": "", "matched_key": "" })).is_empty());
1101 }
1102
1103 /// The toolkit's requests, as `@actions/cache` 4 and `@actions/artifact`
1104 /// 2 send them (protobuf-ts, proto field names, no defaults).
1105 #[test]
1106 fn twirp_requests_read_in_the_toolkits_spelling() {
1107 let create_cache = json!({ "key": "node-cache-Linux-x64-npm-abc", "version": "a7f2c1e0" });
1108 assert_eq!(text(&create_cache, "key"), "node-cache-Linux-x64-npm-abc");
1109 let lookup = json!({ "key": "k", "restore_keys": ["k-", "x-"], "version": "v" });
1110 assert_eq!(field(&lookup, "restore_keys").as_array().unwrap().len(), 2);
1111 let finalize = json!({ "key": "k", "size_bytes": "1048576", "version": "v" });
1112 assert_eq!(number(&finalize, "size_bytes"), Some(1_048_576));
1113 let create_artifact = json!({
1114 "workflow_run_backend_id": "run_1",
1115 "workflow_job_run_backend_id": "job_1",
1116 "name": "dist",
1117 "expires_at": "2026-10-13T00:00:00Z",
1118 "version": 4
1119 });
1120 assert_eq!(text(&create_artifact, "workflow_job_run_backend_id"), "job_1");
1121 let list = json!({ "workflow_run_backend_id": "run_1", "workflow_job_run_backend_id": "job_1", "id_filter": "42", "name_filter": "dist" });
1122 assert_eq!(number(&list, "id_filter"), Some(42));
1123 assert_eq!(text(&list, "name_filter"), "dist");
1124 // A client writing JSON names instead reads the same.
1125 let camel = json!({ "workflowRunBackendId": "run_1", "sizeBytes": 3 });
1126 assert_eq!(text(&camel, "workflow_run_backend_id"), "run_1");
1127 assert_eq!(number(&camel, "size_bytes"), Some(3));
1128 }
1129
1130 #[test]
1131 fn the_runtime_token_names_its_run_and_job() {
1132 let payload = URL_SAFE_NO_PAD.encode(json!({ "job": "job_1", "run": "run_1" }).to_string());
1133 let token = format!("h.{payload}.s");
1134 assert_eq!(runtime_job(&token).as_deref(), Some("job_1"));
1135 assert_eq!(backend_ids(&token), ("run_1".to_owned(), "job_1".to_owned()));
1136 assert_eq!(runtime_job("deadbeef"), None);
1137 }
1138
1139 /// sccache 0.18's storage check, at server start: a lookup of
1140 /// `sccache/.sccache_check`. The actions service answers a miss with
1141 /// `Ok(None)`, `{"ok":true,"value":null}`, which was read back as a
1142 /// malformed outcome, and every lookup that missed was a 500
1143 /// ("Server startup failed: cache storage failed to read").
1144 #[test]
1145 fn sccaches_first_lookup_misses_with_a_204() {
1146 let url = worker::Url::parse(
1147 "https://api.g1t.sh/actions/toolkit/_apis/artifactcache/cache?keys=sccache/.sccache_check&version=sccache-v0.18.0",
1148 )
1149 .unwrap();
1150 assert_eq!(query_in(&url, "keys").as_deref(), Some("sccache/.sccache_check"));
1151 assert_eq!(query_in(&url, "version").as_deref(), Some("sccache-v0.18.0"));
1152
1153 // As the actions service replies (`g1t_kit::reply`), and the API
1154 // reads it (`g1t_kit::call`).
1155 let wire = serde_json::to_string(&Outcome::<Option<CacheHit>>::Ok(None)).unwrap();
1156 assert_eq!(wire, r#"{"ok":true,"value":null}"#);
1157 let found: Outcome<Option<CacheHit>> = g1t_kit::read_answer("cache_lookup", &wire).unwrap();
1158 assert_eq!(lookup_answer(found, "sccache-v0.18.0", "https://api.g1t.sh"), (204, None));
1159
1160 // Once saved, the same lookup is a hit with its download link.
1161 let hit = CacheHit {
1162 key: "sccache/.sccache_check".into(),
1163 object: "c/repo_1/cache_1".into(),
1164 size: 13,
1165 created_at: "2026-10-08T12:00:00.000Z".into(),
1166 blob: Some("tok.sig".into()),
1167 };
1168 let wire = serde_json::to_string(&Outcome::Ok(Some(hit))).unwrap();
1169 let found: Outcome<Option<CacheHit>> = g1t_kit::read_answer("cache_lookup", &wire).unwrap();
1170 let (status, body) = lookup_answer(found, "sccache-v0.18.0", "https://api.g1t.sh");
1171 let body = body.unwrap();
1172 assert_eq!(status, 200);
1173 assert_eq!(body["cacheKey"], "sccache/.sccache_check");
1174 assert_eq!(body["cacheVersion"], "sccache-v0.18.0");
1175 assert_eq!(body["archiveLocation"], "https://api.g1t.sh/actions/toolkit/blobs/tok.sig");
1176
1177 // A refusal keeps its status and says why.
1178 let refused = Outcome::<Option<CacheHit>>::fail(FailureCode::Unauthenticated, "That job is not running.");
1179 let (status, body) = lookup_answer(refused, "v", "https://api.g1t.sh");
1180 assert_eq!((status, body.unwrap()["message"].as_str()), (401, Some("That job is not running.")));
1181
1182 // An answer that does not read names its method and the cause.
1183 let unread = g1t_kit::read_answer::<Outcome<CacheHit>>("cache_lookup", r#"{"ok":true,"value":null}"#).unwrap_err();
1184 assert!(unread.to_string().contains("cache_lookup answered with what could not be read"), "{unread}");
1185 }
1186
1187 #[test]
1188 fn a_job_is_told_where_the_toolkit_s_services_are() {
1189 let vars = runtime_variables("https://api.g1t.sh", "tok", false);
1190 // Twirp paths are resolved against the root of ACTIONS_RESULTS_URL.
1191 assert_eq!(vars["ACTIONS_RESULTS_URL"], "https://api.g1t.sh/");
1192 // The older protocol appends `_apis/artifactcache/…`.
1193 assert_eq!(vars["ACTIONS_CACHE_URL"], "https://api.g1t.sh/actions/toolkit/");
1194 assert_eq!(vars["ACTIONS_CACHE_SERVICE_V2"], "True");
1195 assert!(vars.get("ACTIONS_ID_TOKEN_REQUEST_URL").is_none());
1196 let vars = runtime_variables("https://api.g1t.sh", "tok", true);
1197 // core.getIDToken appends `&audience=…`.
1198 assert_eq!(vars["ACTIONS_ID_TOKEN_REQUEST_URL"], "https://api.g1t.sh/actions/oidc/token?api-version=2.0");
1199 assert_eq!(vars["ACTIONS_ID_TOKEN_REQUEST_TOKEN"], "tok");
1200 }
1201
1202 #[test]
1203 fn artifacts_are_listed_as_the_toolkit_reads_them() {
1204 let artifact = Artifact {
1205 id: 7,
1206 name: "dist".into(),
1207 size: 10,
1208 digest: None,
1209 format: "zip".into(),
1210 run_id: "run_1".into(),
1211 job_id: "job_1".into(),
1212 repo_id: "repo_1".into(),
1213 expired: false,
1214 created_at: "2026-10-08T12:00:00.000Z".into(),
1215 updated_at: "2026-10-08T12:00:00.000Z".into(),
1216 expires_at: "2026-10-22T12:00:00.000Z".into(),
1217 head_branch: None,
1218 head_sha: None,
1219 };
1220 let shown = listed(&artifact);
1221 assert_eq!(shown["database_id"], "7");
1222 assert_eq!(shown["size"], "10");
1223 assert_eq!(shown["workflow_job_run_backend_id"], "job_1");
1224 }
1225}